flagon-io/g1t

public

Git for AI scale: a forge for thousands of agents working on the same code at once.

g1t/services/work/src/reviews.rs

376 lines13,817 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

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