flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/reviews.rs

383 lines14,352 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/// The account reviews by a g1t agent are attributed to. It is not a user
26/// anyone can sign in as.
27pub(crate) const AGENT_ID: &str = "usr_g1t_agent";
28pub(crate) const AGENT_NAME: &str = "g1t-agent";
29
30#[derive(Deserialize)]
31struct OtherRow {
32 number: u32,
33 title: String,
34 issue_number: Option<u32>,
35 files: Option<String>,
36}
37
38#[derive(Deserialize)]
39struct ReviewRunRow {
40 id: String,
41 pull_id: String,
42 token_hash: String,
43 finished_at: Option<String>,
44}
45
46fn hash(token: &str) -> String {
47 hex::encode(Sha256::digest(token.as_bytes()))
48}
49
50impl Work {
51 pub(crate) async fn pull_by_id(&self, id: &str) -> Result<Option<Pull>> {
52 Ok(self
53 .db
54 .prepare(format!("SELECT {PULL_COLUMNS} FROM pulls WHERE id = ?"))
55 .bind(&[id.into()])?
56 .first::<PullRow>(None)
57 .await?
58 .map(Pull::from))
59 }
60
61 /// Works out which files a pull request changes and records them, so
62 /// that overlaps between pull requests can be found without comparing
63 /// every one of them each time.
64 pub(crate) async fn refresh_files(&self, pull: &Pull) -> Result<Vec<ChangedFile>> {
65 // Its author can read both the repository and the pull request's source.
66 let viewer = self.author_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 };
152 self.timing.rpc(g1t_kit::call(&self.repos, "behind", &asked)).await
153 }
154
155 pub(crate) async fn review_pending(&self, pull_id: &str) -> Result<bool> {
156 if let Some(found) = self.prefetched_pull(pull_id) {
157 return Ok(found.first::<crate::rows::ValueRow>(crate::prefetch::Slot::ReviewPending)?.is_some());
158 }
159 Ok(self
160 .db
161 .prepare(
162 "SELECT id AS value FROM review_runs
163 WHERE pull_id = ? AND finished_at IS NULL AND created_at > ? LIMIT 1",
164 )
165 // A sandbox that never reported is forgotten after a while.
166 .bind(&[
167 pull_id.into(),
168 rfc3339(now_ms().saturating_sub(crate::prefetch::REVIEW_PENDING_MS)).into(),
169 ])?
170 .first::<crate::rows::ValueRow>(None)
171 .await?
172 .is_some())
173 }
174
175 /// Begins a review of a pull request by a g1t agent and returns what a
176 /// sandbox needs to carry it out.
177 pub(crate) async fn start_review(&self, a: StartReviewArgs) -> Result<Outcome<ReviewJob>> {
178 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
179 return Ok(Outcome::fail(
180 FailureCode::NotFound,
181 "Pull request not found.",
182 ));
183 };
184 if pull.status != PullStatus::Open {
185 return Ok(Outcome::fail(
186 FailureCode::Conflict,
187 "A pull request can be reviewed once it is ready for review.",
188 ));
189 }
190 let viewer = self.author_viewer(&pull).await?;
191 let repo: Outcome<Repo> = g1t_kit::call(
192 &self.repos,
193 "get_by_id",
194 &GetByIdArgs {
195 id: pull.repo_id.clone(),
196 viewer,
197 },
198 )
199 .await?;
200 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
201 return Ok(Outcome::fail(
202 FailureCode::NotFound,
203 "Pull request not found.",
204 ));
205 };
206 let head: Option<String> = g1t_kit::call(
207 &self.repos,
208 "head",
209 &HeadArgs {
210 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
211 branch: pull
212 .branch
213 .clone()
214 .unwrap_or_else(|| repo.default_branch.clone()),
215 },
216 )
217 .await?;
218 let Some(commit) = head else {
219 return Ok(Outcome::fail(
220 FailureCode::Conflict,
221 "This pull request has no commits to review.",
222 ));
223 };
224 let issue = match pull.issue {
225 Some(number) => self.issue(&pull.repo_id, number).await?,
226 None => None,
227 };
228
229 let now = now_ms();
230 let run_id = new_id("rvw", now);
231 let mut bytes = [0u8; 32];
232 getrandom::getrandom(&mut bytes).expect("no source of randomness");
233 let token = hex::encode(bytes);
234 self.db
235 .prepare(
236 "INSERT INTO review_runs (id, pull_id, token_hash, created_at) VALUES (?, ?, ?, ?)",
237 )
238 .bind(&[
239 run_id.as_str().into(),
240 pull.id.as_str().into(),
241 hash(&token).into(),
242 rfc3339(now).into(),
243 ])?
244 .run()
245 .await?;
246
247 let path = RepoPath {
248 namespace: repo.namespace,
249 name: repo.name,
250 };
251 Ok(Outcome::Ok(ReviewJob {
252 run_id,
253 token,
254 source: pull.fork.clone().unwrap_or_else(|| path.clone()),
255 commit,
256 repo: path,
257 default_branch: repo.default_branch,
258 number: pull.number,
259 title: pull.title,
260 description: pull.body.unwrap_or_default(),
261 issue,
262 author: pull.author,
263 }))
264 }
265
266 /// Records the review a sandbox's agent wrote: its verdict, its summary
267 /// and its comments on lines. The run's token is the only credential.
268 pub(crate) async fn report_review(&self, a: ReportReviewArgs) -> Result<Outcome<bool>> {
269 let run = self
270 .db
271 .prepare("SELECT id, pull_id, token_hash, finished_at FROM review_runs WHERE id = ?")
272 .bind(&[a.run_id.as_str().into()])?
273 .first::<ReviewRunRow>(None)
274 .await?;
275 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
276 return Ok(Outcome::fail(FailureCode::NotFound, "Review not found."));
277 };
278 if run.finished_at.is_some() {
279 return Ok(Outcome::fail(
280 FailureCode::Conflict,
281 "This review has already been reported.",
282 ));
283 }
284 let Some(pull) = self.pull_by_id(&run.pull_id).await? else {
285 return Ok(Outcome::fail(FailureCode::NotFound, "Review not found."));
286 };
287
288 let now = now_ms();
289 let timestamp = rfc3339(now);
290 let signature = a
291 .model
292 .as_deref()
293 .map(|model| format!("\n\n_Reviewed by a g1t agent on {model}._"))
294 .unwrap_or_default();
295 let summary = match &a.error {
296 Some(error) => format!("The review could not be completed: {error}"),
297 None => {
298 let body: String = a.body.trim().chars().take(MAX_REVIEW_CHARS).collect();
299 format!("{body}{signature}")
300 }
301 };
302 // A failed review carries no verdict.
303 let verdict = a.verdict.filter(|_| a.error.is_none());
304 let known: HashSet<&str> = pull.files.iter().map(|file| file.path.as_str()).collect();
305
306 let insert = "INSERT INTO comments
307 (id, repo_id, number, author_id, author_name, body, path, line, verdict, created_at)
308 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)";
309 let mut statements = Vec::new();
310 // Line comments first, so the summary and verdict read as the conclusion.
311 for comment in a
312 .comments
313 .iter()
314 .filter(|comment| a.error.is_none() && !comment.body.trim().is_empty())
315 // A line in a file the pull request does not change is not a review of it.
316 .filter(|comment| known.is_empty() || known.contains(comment.path.as_str()))
317 .take(MAX_REVIEW_COMMENTS)
318 {
319 let body: String = comment.body.trim().chars().take(MAX_REVIEW_CHARS).collect();
320 statements.push(self.db.prepare(insert).bind(&[
321 new_id("cmt", now).into(),
322 pull.repo_id.as_str().into(),
323 pull.number.into(),
324 AGENT_ID.into(),
325 AGENT_NAME.into(),
326 body.into(),
327 comment.path.as_str().into(),
328 optional_number(Some(comment.line).filter(|line| *line > 0)),
329 JsValue::NULL,
330 timestamp.as_str().into(),
331 ])?);
332 }
333 statements.push(self.db.prepare(insert).bind(&[
334 new_id("cmt", now).into(),
335 pull.repo_id.as_str().into(),
336 pull.number.into(),
337 AGENT_ID.into(),
338 AGENT_NAME.into(),
339 summary.into(),
340 JsValue::NULL,
341 JsValue::NULL,
342 verdict.map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
343 timestamp.as_str().into(),
344 ])?);
345 statements.push(
346 self.db
347 .prepare("UPDATE review_runs SET finished_at = ?, verdict = ? WHERE id = ?")
348 .bind(&[
349 timestamp.as_str().into(),
350 verdict.map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
351 run.id.as_str().into(),
352 ])?,
353 );
354 statements.push(
355 self.db
356 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
357 .bind(&[timestamp.as_str().into(), pull.id.as_str().into()])?,
358 );
359 // The review g1t asked for itself is no longer under way.
360 statements.push(
361 self.db
362 .prepare(
363 "UPDATE pulls SET working_on = NULL, working_until = NULL
364 WHERE id = ? AND working_on = 'review'",
365 )
366 .bind(&[pull.id.as_str().into()])?,
367 );
368 self.db.batch(statements).await?;
369 self.publish_as(
370 "review.completed",
371 &pull.repo_id,
372 None,
373 ReviewEvent {
374 pull_id: pull.id.clone(),
375 repo_id: pull.repo_id.clone(),
376 number: pull.number,
377 verdict: verdict.map(Verdict::as_str),
378 },
379 )
380 .await?;
381 Ok(Outcome::Ok(true))
382 }
383}