g1t/services/work/src/reviews.rs

391 lines14,736 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
Thirteen MCP tools and classic token scopes; agents rate their confidence and can be put on an issue in one step18use crate::rows::{PULL_COLUMNS, PullRow};
Agents as a team: lifecycle, merge queue, billing and a new shell19use 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;
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent25/// 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};
Agents as a team: lifecycle, merge queue, billing and a new shell28
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
Thirteen MCP tools and classic token scopes; agents rate their confidence and can be put on an issue in one step53 .prepare(format!("SELECT {PULL_COLUMNS} FROM pulls WHERE id = ?"))
Agents as a team: lifecycle, merge queue, billing and a new shell54 .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>> {
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights64 // 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?;
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();
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily100 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 };
Agents as a team: lifecycle, merge queue, billing and a new shell104 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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily124 /// 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
Agents as a team: lifecycle, merge queue, billing and a new shell139 /// 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 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily145 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
Agents as a team: lifecycle, merge queue, billing and a new shell153 }
154
155 pub(crate) async fn review_pending(&self, pull_id: &str) -> Result<bool> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily156 if let Some(found) = self.prefetched_pull(pull_id) {
157 return Ok(found.first::<crate::rows::ValueRow>(crate::prefetch::Slot::ReviewPending)?.is_some());
158 }
Agents as a team: lifecycle, merge queue, billing and a new shell159 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(),
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily168 rfc3339(now_ms().saturating_sub(crate::prefetch::REVIEW_PENDING_MS)).into(),
Agents as a team: lifecycle, merge queue, billing and a new shell169 ])?
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 }
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights190 let viewer = self.owner_viewer(&pull).await?;
Agents as a team: lifecycle, merge queue, billing and a new shell191 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?;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look200 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
Agents as a team: lifecycle, merge queue, billing and a new shell201 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 };
Auto model routing: the cheapest tier that can do each piece of work, a retry goes up a tier, and each run records its tier251 let mut sensitive: Vec<String> = Vec::new();
252 for kind in pull.files.iter().filter_map(|file| crate::confidence::sensitive(&file.path)) {
253 if !sensitive.iter().any(|seen| seen == kind) {
254 sensitive.push(kind.to_owned());
255 }
256 }
Agents as a team: lifecycle, merge queue, billing and a new shell257 Ok(Outcome::Ok(ReviewJob {
258 run_id,
259 token,
260 source: pull.fork.clone().unwrap_or_else(|| path.clone()),
261 commit,
262 repo: path,
263 default_branch: repo.default_branch,
264 number: pull.number,
265 title: pull.title,
266 description: pull.body.unwrap_or_default(),
267 issue,
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights268 author: pull.requested_by.unwrap_or(pull.author),
Auto model routing: the cheapest tier that can do each piece of work, a retry goes up a tier, and each run records its tier269 files: pull.files,
270 sensitive,
Agents as a team: lifecycle, merge queue, billing and a new shell271 }))
272 }
273
274 /// Records the review a sandbox's agent wrote: its verdict, its summary
275 /// and its comments on lines. The run's token is the only credential.
276 pub(crate) async fn report_review(&self, a: ReportReviewArgs) -> Result<Outcome<bool>> {
277 let run = self
278 .db
279 .prepare("SELECT id, pull_id, token_hash, finished_at FROM review_runs WHERE id = ?")
280 .bind(&[a.run_id.as_str().into()])?
281 .first::<ReviewRunRow>(None)
282 .await?;
283 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
284 return Ok(Outcome::fail(FailureCode::NotFound, "Review not found."));
285 };
286 if run.finished_at.is_some() {
287 return Ok(Outcome::fail(
288 FailureCode::Conflict,
289 "This review has already been reported.",
290 ));
291 }
292 let Some(pull) = self.pull_by_id(&run.pull_id).await? else {
293 return Ok(Outcome::fail(FailureCode::NotFound, "Review not found."));
294 };
295
296 let now = now_ms();
297 let timestamp = rfc3339(now);
298 let signature = a
299 .model
300 .as_deref()
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent301 .map(|model| format!("\n\n_Reviewed by g1t on {model}._"))
Agents as a team: lifecycle, merge queue, billing and a new shell302 .unwrap_or_default();
303 let summary = match &a.error {
304 Some(error) => format!("The review could not be completed: {error}"),
305 None => {
306 let body: String = a.body.trim().chars().take(MAX_REVIEW_CHARS).collect();
307 format!("{body}{signature}")
308 }
309 };
310 // A failed review carries no verdict.
311 let verdict = a.verdict.filter(|_| a.error.is_none());
312 let known: HashSet<&str> = pull.files.iter().map(|file| file.path.as_str()).collect();
313
314 let insert = "INSERT INTO comments
315 (id, repo_id, number, author_id, author_name, body, path, line, verdict, created_at)
316 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)";
317 let mut statements = Vec::new();
318 // Line comments first, so the summary and verdict read as the conclusion.
319 for comment in a
320 .comments
321 .iter()
322 .filter(|comment| a.error.is_none() && !comment.body.trim().is_empty())
323 // A line in a file the pull request does not change is not a review of it.
324 .filter(|comment| known.is_empty() || known.contains(comment.path.as_str()))
325 .take(MAX_REVIEW_COMMENTS)
326 {
327 let body: String = comment.body.trim().chars().take(MAX_REVIEW_CHARS).collect();
328 statements.push(self.db.prepare(insert).bind(&[
329 new_id("cmt", now).into(),
330 pull.repo_id.as_str().into(),
331 pull.number.into(),
332 AGENT_ID.into(),
333 AGENT_NAME.into(),
334 body.into(),
335 comment.path.as_str().into(),
336 optional_number(Some(comment.line).filter(|line| *line > 0)),
337 JsValue::NULL,
338 timestamp.as_str().into(),
339 ])?);
340 }
341 statements.push(self.db.prepare(insert).bind(&[
342 new_id("cmt", now).into(),
343 pull.repo_id.as_str().into(),
344 pull.number.into(),
345 AGENT_ID.into(),
346 AGENT_NAME.into(),
347 summary.into(),
348 JsValue::NULL,
349 JsValue::NULL,
350 verdict.map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
351 timestamp.as_str().into(),
352 ])?);
353 statements.push(
354 self.db
355 .prepare("UPDATE review_runs SET finished_at = ?, verdict = ? WHERE id = ?")
356 .bind(&[
357 timestamp.as_str().into(),
358 verdict.map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
359 run.id.as_str().into(),
360 ])?,
361 );
362 statements.push(
363 self.db
364 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
365 .bind(&[timestamp.as_str().into(), pull.id.as_str().into()])?,
366 );
367 // The review g1t asked for itself is no longer under way.
368 statements.push(
369 self.db
370 .prepare(
371 "UPDATE pulls SET working_on = NULL, working_until = NULL
372 WHERE id = ? AND working_on = 'review'",
373 )
374 .bind(&[pull.id.as_str().into()])?,
375 );
376 self.db.batch(statements).await?;
377 self.publish_as(
378 "review.completed",
379 &pull.repo_id,
380 None,
381 ReviewEvent {
382 pull_id: pull.id.clone(),
383 repo_id: pull.repo_id.clone(),
384 number: pull.number,
385 verdict: verdict.map(Verdict::as_str),
386 },
387 )
388 .await?;
389 Ok(Outcome::Ok(true))
390 }
391}