pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/lifecycle.rs

1,416 lines50,412 bytesCodeBlame
1//! Seeing a pull request through. Once a g1t agent has made a change, g1t
2//! takes each remaining step itself: the acceptance checks, a review by
3//! another agent, sending the author back to address what either found,
4//! and catching up when the branch it would land on has moved. It stops
5//! when the pull request is ready for a person to merge, or when it has
6//! tried and a person has to decide.
7//!
8//! This service decides what the next step is and claims it. The runner
9//! service asks, on every event that could change the answer, and carries
10//! the step out in a sandbox.
11
12use g1t_contracts::events::PullEvent;
13use g1t_contracts::repos::{GetByIdArgs, Repo, RepoPath};
14use g1t_contracts::time::rfc3339;
15use g1t_contracts::work::*;
16use g1t_contracts::{FailureCode, Membership, Outcome, Role, User, Viewer};
17use g1t_kit::now_ms;
18use std::collections::HashMap;
19
20use serde::Deserialize;
21use worker::Result;
22use worker::wasm_bindgen::JsValue;
23
24use crate::Work;
25use crate::reviews::{AGENT_ID, AGENT_NAME};
26use crate::rows::ValueRow;
27
28/// How long a claimed step is waited for before it may be taken again.
29const REVIEW_MINUTES: u64 = 20;
30const REVISION_MINUTES: u64 = 60;
31const CATCH_UP_MINUTES: u64 = 30;
32const MERGE_MINUTES: u64 = 2;
33/// Who a merge made by a repository's settings is attributed to. Not an
34/// account: `g1t` cannot be registered.
35const POLICY_ACTOR_ID: &str = "g1t_policy";
36const POLICY_ACTOR_NAME: &str = "g1t";
37/// How much of a failed check's output the author is shown.
38const MAX_CHECK_OUTPUT_CHARS: usize = 4_000;
39const MANAGED_PAGE: u32 = 200;
40
41/// The part of a pull request's row that tracks its lifecycle.
42#[derive(Deserialize)]
43struct Progress {
44 managed: u8,
45 revisions: u32,
46 revised_at: Option<String>,
47 working_on: Option<String>,
48 working_until: Option<String>,
49 stalled: Option<String>,
50}
51
52/// A merge asked for while the pull request was behind, as stored.
53#[derive(serde::Serialize, Deserialize)]
54#[serde(rename_all = "camelCase")]
55struct LandRequest {
56 actor: User,
57 keep_issue_open: bool,
58}
59
60#[derive(Deserialize)]
61struct LandRow {
62 land_requested: Option<String>,
63 land_requested_at: Option<String>,
64 stalled: Option<String>,
65}
66
67#[derive(Deserialize)]
68struct FinishedReview {
69 finished_at: String,
70 verdict: Option<Verdict>,
71}
72
73#[derive(Deserialize)]
74struct ReviewNote {
75 body: String,
76 path: Option<String>,
77 line: Option<u32>,
78}
79
80/// What the author is being sent back to address.
81pub(crate) enum Feedback {
82 FailedChecks,
83 /// The review that finished at this time.
84 Review(String),
85 /// What a person who asked for changes wrote since the last revision.
86 Person(PersonRequest),
87}
88
89/// A person's request for changes that still stands: their latest verdict,
90/// made after the agent last revised.
91#[derive(Clone, Debug, Deserialize)]
92pub(crate) struct PersonRequest {
93 author_id: String,
94 author_name: String,
95 created_at: String,
96}
97
98/// What should happen next, if it is g1t's turn.
99pub(crate) enum Next {
100 /// A step is under way, or it is a person's turn.
101 Wait,
102 Review,
103 Revise(Feedback),
104 CatchUp,
105 /// Land it, because the repository says ready pull requests land.
106 Merge,
107}
108
109/// Whether g1t made this pull request, and so sees it through.
110pub(crate) fn made_by_g1t(pull: &Pull) -> bool {
111 pull.runtime == Runtime::Hosted && pull.agent == AGENT_NAME && pull.fork.is_some()
112}
113
114fn at(stage: Stage, detail: impl Into<String>, revisions: u32) -> Lifecycle {
115 Lifecycle {
116 stage,
117 detail: detail.into(),
118 revisions,
119 }
120}
121
122fn times(count: u32) -> String {
123 match count {
124 1 => "once".to_owned(),
125 2 => "twice".to_owned(),
126 count => format!("{count} times"),
127 }
128}
129
130/// Everything the next step depends on.
131struct Facts {
132 /// Still being made: not yet marked ready for review.
133 draft: bool,
134 /// Why g1t stopped, if it has.
135 stalled: Option<String>,
136 /// The step under way, if one was claimed and is still being waited for.
137 working_on: Option<String>,
138 check_status: Option<CheckStatus>,
139 /// A review someone asked for is being written.
140 review_pending: bool,
141 /// Whether the issue has acceptance checks at all.
142 has_checks: bool,
143 revisions: u32,
144 /// The latest finished review of the change as it is now.
145 review: Option<FinishedReview>,
146 /// Whether the branch it would land on has moved without it.
147 behind: bool,
148 /// Whether the repository lands a ready pull request by itself.
149 auto_merge: bool,
150 /// Whether the repository refuses to merge one that is behind.
151 require_up_to_date: bool,
152 /// Whether a second agent reviews it without being asked.
153 agent_review: bool,
154 /// How many times the author may be sent back.
155 max_revisions: u32,
156 /// What the repository's approval rule still wants, if anything.
157 approvals_missing: Option<String>,
158 /// A person asked for changes since the agent last revised.
159 person_request: Option<PersonRequest>,
160 /// Its place in the merge queue, and what is ahead of it there.
161 queued: Option<(QueueState, Vec<u32>)>,
162}
163
164/// Where a pull request stands, and the step to take if it is g1t's turn.
165///
166/// The order is: nothing while a step is under way; a person asking for
167/// changes is answered first; the checks must pass; then a review must
168/// approve; then it must be up to date. A person's or an agent's request
169/// for changes, or a failed check, sends the author back, a limited number
170/// of times, after which a person is asked.
171fn decide(facts: Facts) -> (Lifecycle, Next) {
172 let revisions = facts.revisions;
173 let wait = |stage, detail: &str| (at(stage, detail, revisions), Next::Wait);
174 let exhausted = revisions >= facts.max_revisions;
175
176 if facts.draft {
177 return wait(Stage::Working, "A g1t agent is making the change.");
178 }
179 if let Some(reason) = &facts.stalled {
180 return wait(Stage::NeedsYou, reason);
181 }
182 if let Some((state, ahead)) = &facts.queued {
183 let named = ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ");
184 let detail = match (state, ahead.is_empty()) {
185 (QueueState::Testing, true) => "In the merge queue: being tested on the default branch as it is.".to_owned(),
186 (QueueState::Testing, false) => {
187 format!("In the merge queue: being tested together with {named}, ahead of it.")
188 }
189 (QueueState::Passed, true) => "Passed in the merge queue. Landing.".to_owned(),
190 (QueueState::Passed, false) => {
191 format!("Passed in the merge queue together with {named}. It lands once they have.")
192 }
193 _ => "In the merge queue, waiting for its turn to be tested.".to_owned(),
194 };
195 return wait(Stage::Queued, &detail);
196 }
197 match facts.working_on.as_deref() {
198 Some("revision") => {
199 return wait(
200 Stage::Revising,
201 "The agent is addressing what the checks or the review found.",
202 );
203 }
204 Some("catch_up") => {
205 return wait(
206 Stage::CatchingUp,
207 "The agent is merging in the branch this will land on, which has moved.",
208 );
209 }
210 Some("merge") => return wait(Stage::Ready, "Merging."),
211 Some(_) => return wait(Stage::Reviewing, "A g1t agent is reviewing the change."),
212 None => {}
213 }
214 if matches!(
215 facts.check_status,
216 Some(CheckStatus::Queued | CheckStatus::Running)
217 ) {
218 return wait(Stage::Checking, "The acceptance checks are running.");
219 }
220 if facts.review_pending {
221 return wait(Stage::Reviewing, "A g1t agent is reviewing the change.");
222 }
223 // A person asked for changes: the agent makes them, as it would for a
224 // review it asked for, before anything else.
225 if let Some(request) = &facts.person_request {
226 if exhausted {
227 return wait(
228 Stage::NeedsYou,
229 &format!(
230 "{} asked for changes, and the agent has already revised {}.",
231 request.author_name,
232 times(revisions)
233 ),
234 );
235 }
236 return (
237 at(
238 Stage::Revising,
239 format!(
240 "{} asked for changes. The agent is being sent back to make them.",
241 request.author_name
242 ),
243 revisions,
244 ),
245 Next::Revise(Feedback::Person(request.clone())),
246 );
247 }
248
249 if facts.has_checks {
250 match facts.check_status {
251 Some(CheckStatus::Passed) => {}
252 Some(CheckStatus::Failed) if exhausted => {
253 return wait(
254 Stage::NeedsYou,
255 &format!(
256 "The acceptance checks still fail after the agent revised {}.",
257 times(revisions)
258 ),
259 );
260 }
261 Some(CheckStatus::Failed) => {
262 return (
263 at(
264 Stage::Revising,
265 "The acceptance checks failed. The agent is being sent back to fix them.",
266 revisions,
267 ),
268 Next::Revise(Feedback::FailedChecks),
269 );
270 }
271 Some(CheckStatus::Errored) => {
272 return wait(Stage::NeedsYou, "The acceptance checks could not be run.");
273 }
274 _ => {
275 return wait(
276 Stage::Checking,
277 "Waiting for the acceptance checks to start.",
278 );
279 }
280 }
281 }
282
283 // A second agent reviews it, unless the repository leaves review to people.
284 if facts.agent_review {
285 match facts.review {
286 None => {
287 return (
288 at(
289 Stage::Reviewing,
290 "A g1t agent is about to review the change.",
291 revisions,
292 ),
293 Next::Review,
294 );
295 }
296 Some(FinishedReview { verdict: None, .. }) => {
297 return wait(Stage::NeedsYou, "The review could not be completed.");
298 }
299 Some(FinishedReview {
300 verdict: Some(Verdict::RequestChanges),
301 ..
302 }) if exhausted => {
303 return wait(
304 Stage::NeedsYou,
305 &format!(
306 "The review still asks for changes after the agent revised {}.",
307 times(revisions)
308 ),
309 );
310 }
311 Some(FinishedReview {
312 verdict: Some(Verdict::RequestChanges),
313 finished_at,
314 }) => {
315 return (
316 at(
317 Stage::Revising,
318 "The review asked for changes. The agent is being sent back to make them.",
319 revisions,
320 ),
321 Next::Revise(Feedback::Review(finished_at)),
322 );
323 }
324 Some(FinishedReview {
325 verdict: Some(Verdict::Approve),
326 ..
327 }) => {}
328 }
329 }
330
331 // Only where the repository insists is catching up a step of its own,
332 // followed by the checks again. Elsewhere it happens as part of merging.
333 if facts.behind && facts.require_up_to_date {
334 return (
335 at(
336 Stage::CatchingUp,
337 "The branch it will land on has moved. The agent is catching up.",
338 revisions,
339 ),
340 Next::CatchUp,
341 );
342 }
343 // The repository wants approvals this does not have yet: people's turn,
344 // so it is shown as needing someone, not as g1t still working.
345 if let Some(missing) = &facts.approvals_missing {
346 return wait(Stage::NeedsYou, missing);
347 }
348 if facts.auto_merge {
349 return (
350 at(
351 Stage::Ready,
352 "Everything this repository asks for is met. Merging, as its settings say.",
353 revisions,
354 ),
355 Next::Merge,
356 );
357 }
358 wait(
359 Stage::Ready,
360 if facts.behind {
361 "Ready to merge. Merging brings it up to date with the default branch first."
362 } else {
363 "Everything this repository asks for is met. Ready to merge."
364 },
365 )
366}
367
368impl Work {
369 /// Where a pull request stands and what g1t does next, remembered so
370 /// lists can show it without working it out again. `None` for one g1t
371 /// is not seeing through.
372 pub(crate) async fn assess(
373 &self,
374 pull: &Pull,
375 issue: &Option<Issue>,
376 behind: bool,
377 ) -> Result<Option<(Lifecycle, Next)>> {
378 let assessed = self.assess_now(pull, issue, behind).await?;
379 if let Some((lifecycle, _)) = &assessed {
380 self.remember(&pull.id, lifecycle).await?;
381 }
382 Ok(assessed)
383 }
384
385 /// Saves where a pull request stands, for [`Self::remembered`].
386 pub(crate) async fn remember(&self, pull_id: &str, lifecycle: &Lifecycle) -> Result<()> {
387 let stage = serde_json::to_value(lifecycle.stage)?;
388 self.db
389 .prepare("UPDATE pulls SET stage = ?, stage_detail = ? WHERE id = ?")
390 .bind(&[
391 stage.as_str().unwrap_or_default().into(),
392 lifecycle.detail.as_str().into(),
393 pull_id.into(),
394 ])?
395 .run()
396 .await?;
397 Ok(())
398 }
399
400 async fn assess_now(
401 &self,
402 pull: &Pull,
403 issue: &Option<Issue>,
404 behind: bool,
405 ) -> Result<Option<(Lifecycle, Next)>> {
406 if !pull.status.is_active() {
407 return Ok(None);
408 }
409 let Some(progress) = self
410 .db
411 .prepare(
412 "SELECT managed, revisions, revised_at, working_on, working_until, stalled
413 FROM pulls WHERE id = ?",
414 )
415 .bind(&[pull.id.as_str().into()])?
416 .first::<Progress>(None)
417 .await?
418 .filter(|progress| progress.managed != 0)
419 else {
420 return Ok(None);
421 };
422 let now = rfc3339(now_ms());
423 let working_on = progress
424 .working_until
425 .as_deref()
426 .is_some_and(|until| until > now.as_str())
427 .then(|| progress.working_on.clone().unwrap_or_default());
428 let review = self
429 .db
430 .prepare(
431 "SELECT finished_at, verdict FROM review_runs
432 WHERE pull_id = ? AND finished_at IS NOT NULL ORDER BY id DESC LIMIT 1",
433 )
434 .bind(&[pull.id.as_str().into()])?
435 .first::<FinishedReview>(None)
436 .await?
437 // A review of what the change was before its last revision says
438 // nothing about what it is now.
439 .filter(|review| {
440 progress
441 .revised_at
442 .as_deref()
443 .is_none_or(|revised| review.finished_at.as_str() >= revised)
444 });
445 let settings = self.settings(&pull.repo_id).await?;
446 Ok(Some(decide(Facts {
447 draft: pull.status == PullStatus::Draft,
448 stalled: progress.stalled,
449 working_on,
450 check_status: pull.check_status,
451 review_pending: self.review_pending(&pull.id).await?,
452 has_checks: issue.as_ref().is_some_and(|issue| !issue.checks.is_empty()),
453 revisions: progress.revisions,
454 review,
455 behind,
456 auto_merge: settings.auto_merge,
457 require_up_to_date: settings.require_up_to_date,
458 agent_review: settings.agent_review,
459 max_revisions: settings.max_revisions,
460 approvals_missing: self.approvals_gap(&settings, pull).await?,
461 person_request: self
462 .person_request(pull, progress.revised_at.as_deref())
463 .await?,
464 queued: self.queued_entry(&pull.id).await?,
465 })))
466 }
467
468 /// The latest request for changes by a person other than the author,
469 /// if it is that person's latest verdict and came after the last
470 /// revision.
471 async fn person_request(
472 &self,
473 pull: &Pull,
474 revised_at: Option<&str>,
475 ) -> Result<Option<PersonRequest>> {
476 #[derive(Deserialize)]
477 struct Verdicts {
478 author_id: String,
479 author_name: String,
480 verdict: Verdict,
481 created_at: String,
482 }
483 let rows = self
484 .db
485 .prepare(
486 "SELECT author_id, author_name, verdict, created_at FROM comments
487 WHERE repo_id = ? AND number = ? AND verdict IS NOT NULL
488 AND author_id != ? AND author_id != ?
489 ORDER BY id",
490 )
491 .bind(&[
492 pull.repo_id.as_str().into(),
493 pull.number.into(),
494 pull.author.id.as_str().into(),
495 AGENT_ID.into(),
496 ])?
497 .all()
498 .await?
499 .results::<Verdicts>()?;
500 // Each person's latest verdict is the one that stands.
501 let mut latest: HashMap<String, Verdicts> = HashMap::new();
502 for row in rows {
503 latest.insert(row.author_id.clone(), row);
504 }
505 Ok(latest
506 .into_values()
507 .filter(|row| row.verdict == Verdict::RequestChanges)
508 .filter(|row| revised_at.is_none_or(|revised| row.created_at.as_str() > revised))
509 .max_by(|a, b| a.created_at.cmp(&b.created_at))
510 .map(|row| PersonRequest {
511 author_id: row.author_id,
512 author_name: row.author_name,
513 created_at: row.created_at,
514 }))
515 }
516
517 /// Marks a pull request a g1t agent has just opened as one g1t sees
518 /// through.
519 pub(crate) async fn manage(&self, pull: &Pull) -> Result<()> {
520 if !made_by_g1t(pull) {
521 return Ok(());
522 }
523 self.db
524 .prepare("UPDATE pulls SET managed = 1 WHERE id = ?")
525 .bind(&[pull.id.as_str().into()])?
526 .run()
527 .await?;
528 Ok(())
529 }
530
531 /// Takes a step for a pull request, if nobody else has. One statement,
532 /// so that two callers cannot both take it.
533 async fn claim(&self, pull_id: &str, step: &str, minutes: u64, revising: bool) -> Result<bool> {
534 let now = now_ms();
535 let revision = if revising {
536 ", revisions = revisions + 1, revised_at = ?1"
537 } else {
538 ""
539 };
540 Ok(self
541 .db
542 .prepare(format!(
543 "UPDATE pulls SET working_on = ?2, working_until = ?3{revision}
544 WHERE id = ?4 AND status = 'open' AND stalled IS NULL
545 AND (working_until IS NULL OR working_until < ?1)
546 RETURNING id AS value"
547 ))
548 .bind(&[
549 rfc3339(now).into(),
550 step.into(),
551 rfc3339(now + minutes * 60 * 1000).into(),
552 pull_id.into(),
553 ])?
554 .first::<ValueRow>(None)
555 .await?
556 .is_some())
557 }
558
559 /// What the author is told when sent back: the checks that failed and
560 /// what they printed, or the review and its comments on lines.
561 async fn feedback(&self, pull: &Pull, feedback: &Feedback) -> Result<String> {
562 match feedback {
563 Feedback::FailedChecks => {
564 let failed: Vec<String> = self
565 .latest_checks(&pull.id)
566 .await?
567 .map(|run| run.results)
568 .unwrap_or_default()
569 .into_iter()
570 .filter(|result| !result.passed)
571 .map(|result| {
572 let length = result.output.chars().count();
573 let output: String = result
574 .output
575 .chars()
576 .skip(length.saturating_sub(MAX_CHECK_OUTPUT_CHARS))
577 .collect();
578 let exit = result
579 .exit_code
580 .map_or("it was stopped for taking too long".to_owned(), |code| {
581 format!("exit code {code}")
582 });
583 format!("`{}` failed ({exit}):\n\n{}", result.command, output.trim())
584 })
585 .collect();
586 Ok(format!(
587 "These acceptance checks were run against your change in a clean sandbox and failed.\n\n{}",
588 failed.join("\n\n")
589 ))
590 }
591 Feedback::Review(finished_at) => {
592 // Everything a review says is recorded at the moment it finished.
593 let notes = self
594 .db
595 .prepare(
596 "SELECT body, path, line FROM comments
597 WHERE repo_id = ? AND number = ? AND author_id = ? AND created_at = ?
598 ORDER BY id",
599 )
600 .bind(&[
601 pull.repo_id.as_str().into(),
602 pull.number.into(),
603 AGENT_ID.into(),
604 finished_at.as_str().into(),
605 ])?
606 .all()
607 .await?
608 .results::<ReviewNote>()?;
609 let mut on_lines = Vec::new();
610 let mut summary = String::new();
611 for note in notes {
612 match (note.path, note.line) {
613 (Some(path), Some(line)) => {
614 on_lines.push(format!("- `{path}` line {line}: {}", note.body));
615 }
616 (Some(path), None) => on_lines.push(format!("- `{path}`: {}", note.body)),
617 (None, _) => summary = note.body,
618 }
619 }
620 let mut text = format!(
621 "Another agent reviewed your change and asked for changes.\n\n{summary}"
622 );
623 if !on_lines.is_empty() {
624 text.push_str("\n\nIts comments on lines:\n");
625 text.push_str(&on_lines.join("\n"));
626 }
627 Ok(text)
628 }
629 Feedback::Person(request) => {
630 // What they wrote since the agent last revised, which their
631 // request for changes closes.
632 let revised: Option<String> = self
633 .db
634 .prepare("SELECT revised_at AS value FROM pulls WHERE id = ?")
635 .bind(&[pull.id.as_str().into()])?
636 .first::<Option<String>>(Some("value"))
637 .await?
638 .flatten();
639 let notes = self
640 .db
641 .prepare(
642 "SELECT body, path, line FROM comments
643 WHERE repo_id = ? AND number = ? AND author_id = ?
644 AND created_at > ? AND created_at <= ?
645 ORDER BY id",
646 )
647 .bind(&[
648 pull.repo_id.as_str().into(),
649 pull.number.into(),
650 request.author_id.as_str().into(),
651 revised.unwrap_or_default().into(),
652 request.created_at.as_str().into(),
653 ])?
654 .all()
655 .await?
656 .results::<ReviewNote>()?;
657 let mut on_lines = Vec::new();
658 let mut said = Vec::new();
659 for note in notes {
660 match (note.path, note.line) {
661 (Some(path), Some(line)) => {
662 on_lines.push(format!("- `{path}` line {line}: {}", note.body));
663 }
664 (Some(path), None) => on_lines.push(format!("- `{path}`: {}", note.body)),
665 (None, _) => said.push(note.body),
666 }
667 }
668 let mut text = format!(
669 "{} reviewed your change and asked for changes.\n\n{}",
670 request.author_name,
671 said.join("\n\n")
672 );
673 if !on_lines.is_empty() {
674 text.push_str("\n\nTheir comments on lines:\n");
675 text.push_str(&on_lines.join("\n"));
676 }
677 Ok(text)
678 }
679 }
680 }
681
682 pub(crate) async fn advance(&self, a: AdvanceArgs) -> Result<Advance> {
683 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
684 return Ok(Advance::None);
685 };
686 if pull.status != PullStatus::Open {
687 return Ok(Advance::None);
688 }
689 // Its author can read both the repository and the fork.
690 let viewer: Viewer = Some(pull.author.clone());
691 let repo: Outcome<Repo> = g1t_kit::call(
692 &self.repos,
693 "get_by_id",
694 &GetByIdArgs {
695 id: pull.repo_id.clone(),
696 viewer,
697 },
698 )
699 .await?;
700 let (Outcome::Ok(repo), Some(source)) = (repo, pull.fork.clone()) else {
701 return Ok(Advance::None);
702 };
703 let issue = match pull.issue {
704 Some(number) => self.issue(&pull.repo_id, number).await?,
705 None => None,
706 };
707 let behind = self.is_behind(&repo.id, &pull).await?;
708 let Some((lifecycle, next)) = self.assess(&pull, &issue, behind).await? else {
709 return Ok(Advance::None);
710 };
711
712 if matches!(next, Next::Merge) {
713 self.merge_by_policy(&repo, &pull).await?;
714 return Ok(Advance::None);
715 }
716 let (step, minutes) = match &next {
717 Next::Wait | Next::Merge => return Ok(Advance::None),
718 Next::Review => ("review", REVIEW_MINUTES),
719 Next::Revise(_) => ("revision", REVISION_MINUTES),
720 Next::CatchUp => ("catch_up", CATCH_UP_MINUTES),
721 };
722 let feedback = match &next {
723 Next::Revise(feedback) => self.feedback(&pull, feedback).await?,
724 _ => String::new(),
725 };
726 if !self
727 .claim(&pull.id, step, minutes, matches!(next, Next::Revise(_)))
728 .await?
729 {
730 return Ok(Advance::None);
731 }
732 // Said in the conversation, so nobody has to wonder why a review or
733 // a new commit appeared.
734 let told = match &next {
735 Next::Review => {
736 self.db
737 .prepare(
738 "UPDATE pulls SET reviewers = json_insert(reviewers, '$[#]', ?1)
739 WHERE id = ?2 AND NOT EXISTS (
740 SELECT 1 FROM json_each(pulls.reviewers) WHERE json_each.value = ?1)",
741 )
742 .bind(&[AGENT_NAME.into(), pull.id.as_str().into()])?
743 .run()
744 .await?;
745 "requested a review from g1t-agent".to_owned()
746 }
747 Next::Revise(Feedback::FailedChecks) => {
748 "sent g1t-agent back to fix the failed checks".to_owned()
749 }
750 Next::Revise(Feedback::Review(_)) => {
751 "sent g1t-agent back to address the review".to_owned()
752 }
753 _ => format!(
754 "asked g1t-agent to bring this up to date with {}",
755 repo.default_branch
756 ),
757 };
758 self.note(
759 &pull.repo_id,
760 pull.number,
761 (POLICY_ACTOR_ID, POLICY_ACTOR_NAME),
762 &told,
763 )
764 .await?;
765 let job = LifecycleJob {
766 pull_id: pull.id,
767 repo: RepoPath {
768 namespace: repo.namespace,
769 name: repo.name,
770 },
771 number: pull.number,
772 author: pull.author,
773 source,
774 branch: None,
775 default_branch: repo.default_branch,
776 title: pull.title,
777 description: pull.body.unwrap_or_default(),
778 issue,
779 feedback,
780 round: lifecycle.revisions + 1,
781 };
782 Ok(match next {
783 Next::Review => Advance::Review { job },
784 Next::Revise(_) => Advance::Revise { job },
785 Next::CatchUp => Advance::CatchUp { job },
786 Next::Wait | Next::Merge => Advance::None,
787 })
788 }
789
790 /// Lands a pull request that is ready, on the authority of the
791 /// repository's settings instead of a person's click.
792 async fn merge_by_policy(&self, repo: &Repo, pull: &Pull) -> Result<()> {
793 if !self.claim(&pull.id, "merge", MERGE_MINUTES, false).await? {
794 return Ok(());
795 }
796 // g1t acts for the workspace whose members turned this on.
797 let actor = User {
798 id: POLICY_ACTOR_ID.to_owned(),
799 username: POLICY_ACTOR_NAME.to_owned(),
800 verified: true,
801 workspaces: vec![Membership {
802 slug: repo.namespace.clone(),
803 role: Role::Member,
804 }],
805 ..User::default()
806 };
807 let merged = self
808 .merge_pull(PullActionArgs {
809 actor,
810 repo: RepoPath {
811 namespace: repo.namespace.clone(),
812 name: repo.name.clone(),
813 },
814 number: pull.number,
815 summary: String::new(),
816 keep_issue_open: false,
817 ignore_checks: false,
818 })
819 .await?;
820 match merged {
821 Outcome::Ok(_) => Ok(()),
822 // Most likely the branch moved in the moment between: let go, and
823 // the next look at it will catch up and try again.
824 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
825 self.db
826 .prepare(
827 "UPDATE pulls SET working_on = NULL, working_until = NULL
828 WHERE id = ? AND working_on = 'merge'",
829 )
830 .bind(&[pull.id.as_str().into()])?
831 .run()
832 .await?;
833 Ok(())
834 }
835 Outcome::Fail(failure) => {
836 self.stall(StallArgs {
837 pull_id: pull.id.clone(),
838 reason: format!("g1t could not merge this: {}", failure.message),
839 })
840 .await?;
841 Ok(())
842 }
843 }
844 }
845
846 /// Records that a merge was asked for while the pull request was
847 /// behind, and announces it so that the runner brings it up to date.
848 pub(crate) async fn request_landing(
849 &self,
850 pull: &Pull,
851 actor: &User,
852 keep_issue_open: bool,
853 ) -> Result<()> {
854 let now = now_ms();
855 let request = serde_json::to_string(&LandRequest {
856 actor: actor.clone(),
857 keep_issue_open,
858 })?;
859 self.db
860 .prepare(
861 "UPDATE pulls
862 SET land_requested = ?, land_requested_at = ?, stalled = NULL,
863 working_on = 'catch_up', working_until = ?
864 WHERE id = ?",
865 )
866 .bind(&[
867 request.into(),
868 rfc3339(now).into(),
869 rfc3339(now + CATCH_UP_MINUTES * 60 * 1000).into(),
870 pull.id.as_str().into(),
871 ])?
872 .run()
873 .await?;
874 self.publish(
875 "pull.merge_requested",
876 &pull.repo_id,
877 actor,
878 PullEvent {
879 pull_id: pull.id.clone(),
880 repo_id: pull.repo_id.clone(),
881 number: pull.number,
882 issue: pull.issue,
883 ..PullEvent::default()
884 },
885 )
886 .await
887 }
888
889 /// The merge waiting on a pull request, if one was asked for recently
890 /// enough to still stand.
891 async fn land_request(&self, pull_id: &str) -> Result<Option<LandRequest>> {
892 let row = self
893 .db
894 .prepare("SELECT land_requested, land_requested_at, stalled FROM pulls WHERE id = ?")
895 .bind(&[pull_id.into()])?
896 .first::<LandRow>(None)
897 .await?;
898 let oldest = rfc3339(now_ms().saturating_sub(CATCH_UP_MINUTES * 60 * 1000));
899 Ok(row
900 .filter(|row| {
901 row.land_requested_at
902 .as_deref()
903 .is_some_and(|at| at >= oldest.as_str())
904 })
905 .and_then(|row| row.land_requested)
906 .and_then(|request| serde_json::from_str(&request).ok()))
907 }
908
909 /// Whether a merge is waiting on a pull request, and why g1t stopped
910 /// working on it if it did.
911 pub(crate) async fn landing_state(&self, pull_id: &str) -> Result<(bool, Option<String>)> {
912 let stalled = self
913 .db
914 .prepare("SELECT land_requested, land_requested_at, stalled FROM pulls WHERE id = ?")
915 .bind(&[pull_id.into()])?
916 .first::<LandRow>(None)
917 .await?
918 .and_then(|row| row.stalled);
919 Ok((self.land_request(pull_id).await?.is_some(), stalled))
920 }
921
922 async fn forget_landing(&self, pull_id: &str) -> Result<()> {
923 self.db
924 .prepare(
925 "UPDATE pulls SET land_requested = NULL, land_requested_at = NULL WHERE id = ?",
926 )
927 .bind(&[pull_id.into()])?
928 .run()
929 .await?;
930 Ok(())
931 }
932
933 /// Lands a pull request whose head has just moved, if a merge of it was
934 /// waiting for exactly that. What it was caught up to was already
935 /// checked and reviewed apart from the merge, so the checks are not
936 /// waited for again; a repository that wants them rerun turns on
937 /// "require up to date", and then nothing is landed this way.
938 pub(crate) async fn land_if_requested(&self, pull_id: &str) -> Result<()> {
939 let Some(request) = self.land_request(pull_id).await? else {
940 return Ok(());
941 };
942 self.forget_landing(pull_id).await?;
943 let Some(pull) = self.pull_by_id(pull_id).await? else {
944 return Ok(());
945 };
946 let repo: Outcome<Repo> = g1t_kit::call(
947 &self.repos,
948 "get_by_id",
949 &GetByIdArgs {
950 id: pull.repo_id.clone(),
951 viewer: Some(request.actor.clone()),
952 },
953 )
954 .await?;
955 let Outcome::Ok(repo) = repo else {
956 return Ok(());
957 };
958 let merged = self
959 .merge_pull(PullActionArgs {
960 actor: request.actor,
961 repo: RepoPath {
962 namespace: repo.namespace,
963 name: repo.name,
964 },
965 number: pull.number,
966 summary: String::new(),
967 keep_issue_open: request.keep_issue_open,
968 ignore_checks: true,
969 })
970 .await?;
971 if let Outcome::Fail(failure) = merged {
972 self.stall(StallArgs {
973 pull_id: pull.id,
974 reason: format!(
975 "It was brought up to date but could not be merged: {}",
976 failure.message
977 ),
978 })
979 .await?;
980 }
981 Ok(())
982 }
983
984 /// What the runner needs to bring a pull request up to date for a merge
985 /// that is waiting on it.
986 pub(crate) async fn catch_up_job(&self, a: CatchUpJobArgs) -> Result<Option<LifecycleJob>> {
987 if self.land_request(&a.pull_id).await?.is_none() {
988 return Ok(None);
989 }
990 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
991 return Ok(None);
992 };
993 // Its author can read both the repository and the pull request's source.
994 let repo: Outcome<Repo> = g1t_kit::call(
995 &self.repos,
996 "get_by_id",
997 &GetByIdArgs {
998 id: pull.repo_id.clone(),
999 viewer: Some(pull.author.clone()),
1000 },
1001 )
1002 .await?;
1003 let Outcome::Ok(repo) = repo else {
1004 return Ok(None);
1005 };
1006 let path = RepoPath {
1007 namespace: repo.namespace,
1008 name: repo.name,
1009 };
1010 let issue = match pull.issue {
1011 Some(number) => self.issue(&pull.repo_id, number).await?,
1012 None => None,
1013 };
1014 Ok(Some(LifecycleJob {
1015 pull_id: pull.id,
1016 source: pull.fork.unwrap_or_else(|| path.clone()),
1017 repo: path,
1018 number: pull.number,
1019 author: pull.author,
1020 branch: pull.branch,
1021 default_branch: repo.default_branch,
1022 title: pull.title,
1023 description: pull.body.unwrap_or_default(),
1024 issue,
1025 feedback: String::new(),
1026 round: 0,
1027 }))
1028 }
1029
1030 pub(crate) async fn stall(&self, a: StallArgs) -> Result<bool> {
1031 self.db
1032 .prepare(
1033 "UPDATE pulls
1034 SET stalled = ?, working_on = NULL, working_until = NULL,
1035 land_requested = NULL, land_requested_at = NULL
1036 WHERE id = ? AND status = 'open'",
1037 )
1038 .bind(&[a.reason.trim().into(), a.pull_id.as_str().into()])?
1039 .run()
1040 .await?;
1041 self.db
1042 .prepare("UPDATE pulls SET stage = 'needs_you', stage_detail = ? WHERE id = ?")
1043 .bind(&[a.reason.trim().into(), a.pull_id.into()])?
1044 .run()
1045 .await?;
1046 Ok(true)
1047 }
1048
1049 pub(crate) async fn managed_pulls(&self, a: ManagedPullsArgs) -> Result<Vec<String>> {
1050 let rows = self
1051 .db
1052 .prepare(
1053 "SELECT id AS value FROM pulls
1054 WHERE status = 'open' AND managed = 1 AND stalled IS NULL
1055 AND (?1 IS NULL OR repo_id = ?1)
1056 ORDER BY updated_at DESC LIMIT ?2",
1057 )
1058 .bind(&[
1059 a.repo_id.map_or(JsValue::NULL, JsValue::from),
1060 MANAGED_PAGE.into(),
1061 ])?
1062 .all()
1063 .await?
1064 .results::<ValueRow>()?;
1065 Ok(rows.into_iter().map(|row| row.value).collect())
1066 }
1067}
1068
1069#[cfg(test)]
1070mod tests {
1071 use super::*;
1072
1073 const MAX_REVISIONS: u32 = 2;
1074
1075 #[test]
1076 fn approvals_the_repository_wants_are_waited_for() {
1077 let short = || Facts {
1078 review: reviewed(Some(Verdict::Approve)),
1079 approvals_missing: Some("This repository requires 1 approving review.".to_owned()),
1080 ..facts()
1081 };
1082 let (lifecycle, next) = decide(short());
1083 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1084 assert_eq!(
1085 lifecycle.detail,
1086 "This repository requires 1 approving review."
1087 );
1088 assert!(matches!(next, Next::Wait));
1089 // Not even a repository that merges by itself merges without them.
1090 let automatic = Facts {
1091 auto_merge: true,
1092 ..short()
1093 };
1094 assert_eq!(outcome(automatic), (Stage::NeedsYou, "wait"));
1095 }
1096
1097 #[test]
1098 fn a_repository_can_leave_review_to_people() {
1099 let unreviewed = Facts {
1100 agent_review: false,
1101 ..facts()
1102 };
1103 assert_eq!(outcome(unreviewed), (Stage::Ready, "wait"));
1104 let failing = Facts {
1105 agent_review: false,
1106 check_status: Some(CheckStatus::Failed),
1107 ..facts()
1108 };
1109 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1110 }
1111
1112 #[test]
1113 fn a_repository_sets_how_often_the_author_is_sent_back() {
1114 let never = Facts {
1115 max_revisions: 0,
1116 check_status: Some(CheckStatus::Failed),
1117 ..facts()
1118 };
1119 assert_eq!(outcome(never), (Stage::NeedsYou, "wait"));
1120 }
1121
1122 /// A pull request that is ready for review, with checks that passed
1123 /// and nothing else yet.
1124 fn facts() -> Facts {
1125 Facts {
1126 draft: false,
1127 stalled: None,
1128 working_on: None,
1129 check_status: Some(CheckStatus::Passed),
1130 review_pending: false,
1131 has_checks: true,
1132 revisions: 0,
1133 review: None,
1134 behind: false,
1135 auto_merge: false,
1136 require_up_to_date: false,
1137 agent_review: true,
1138 max_revisions: MAX_REVISIONS,
1139 approvals_missing: None,
1140 person_request: None,
1141 queued: None,
1142 }
1143 }
1144
1145 #[test]
1146 fn a_queued_pull_request_waits_in_the_queue() {
1147 let queued = Facts {
1148 review: reviewed(Some(Verdict::Approve)),
1149 queued: Some((QueueState::Testing, vec![12, 14])),
1150 auto_merge: true,
1151 ..facts()
1152 };
1153 let (lifecycle, next) = decide(queued);
1154 assert_eq!(lifecycle.stage, Stage::Queued);
1155 assert!(lifecycle.detail.contains("#12, #14"));
1156 assert!(matches!(next, Next::Wait));
1157 }
1158
1159 fn asked_by_a_person() -> Option<PersonRequest> {
1160 Some(PersonRequest {
1161 author_id: "usr_reviewer".to_owned(),
1162 author_name: "g1t-reviewer".to_owned(),
1163 created_at: "2026-10-02T11:00:00.000Z".to_owned(),
1164 })
1165 }
1166
1167 #[test]
1168 fn a_person_asking_for_changes_sends_the_agent_back() {
1169 let asked = Facts {
1170 review: reviewed(Some(Verdict::Approve)),
1171 approvals_missing: Some("A reviewer has asked for changes.".to_owned()),
1172 person_request: asked_by_a_person(),
1173 ..facts()
1174 };
1175 let (lifecycle, _) = decide(Facts {
1176 person_request: asked_by_a_person(),
1177 ..facts()
1178 });
1179 assert!(lifecycle.detail.starts_with("g1t-reviewer asked for changes"));
1180 assert_eq!(outcome(asked), (Stage::Revising, "revise for a person"));
1181 }
1182
1183 #[test]
1184 fn a_person_is_asked_once_the_revisions_run_out() {
1185 let exhausted = Facts {
1186 person_request: asked_by_a_person(),
1187 revisions: MAX_REVISIONS,
1188 ..facts()
1189 };
1190 assert_eq!(outcome(exhausted), (Stage::NeedsYou, "wait"));
1191 }
1192
1193 fn reviewed(verdict: Option<Verdict>) -> Option<FinishedReview> {
1194 Some(FinishedReview {
1195 finished_at: "2026-10-02T10:00:00.000Z".to_owned(),
1196 verdict,
1197 })
1198 }
1199
1200 /// The stage, and a word for the step to take.
1201 fn outcome(facts: Facts) -> (Stage, &'static str) {
1202 let (lifecycle, next) = decide(facts);
1203 let step = match next {
1204 Next::Wait => "wait",
1205 Next::Review => "review",
1206 Next::Revise(Feedback::FailedChecks) => "revise for checks",
1207 Next::Revise(Feedback::Review(_)) => "revise for review",
1208 Next::Revise(Feedback::Person(_)) => "revise for a person",
1209 Next::CatchUp => "catch up",
1210 Next::Merge => "merge",
1211 };
1212 (lifecycle.stage, step)
1213 }
1214
1215 #[test]
1216 fn nothing_is_started_while_the_agent_is_still_working() {
1217 let draft = Facts {
1218 draft: true,
1219 check_status: None,
1220 ..facts()
1221 };
1222 assert_eq!(outcome(draft), (Stage::Working, "wait"));
1223 }
1224
1225 #[test]
1226 fn checks_come_before_review() {
1227 let unchecked = Facts {
1228 check_status: None,
1229 ..facts()
1230 };
1231 assert_eq!(outcome(unchecked), (Stage::Checking, "wait"));
1232 let running = Facts {
1233 check_status: Some(CheckStatus::Running),
1234 ..facts()
1235 };
1236 assert_eq!(outcome(running), (Stage::Checking, "wait"));
1237 assert_eq!(outcome(facts()), (Stage::Reviewing, "review"));
1238 }
1239
1240 #[test]
1241 fn an_issue_without_checks_goes_straight_to_review() {
1242 let unchecked = Facts {
1243 has_checks: false,
1244 check_status: None,
1245 ..facts()
1246 };
1247 assert_eq!(outcome(unchecked), (Stage::Reviewing, "review"));
1248 }
1249
1250 #[test]
1251 fn failed_checks_send_the_author_back() {
1252 let failed = Facts {
1253 check_status: Some(CheckStatus::Failed),
1254 ..facts()
1255 };
1256 assert_eq!(outcome(failed), (Stage::Revising, "revise for checks"));
1257 }
1258
1259 #[test]
1260 fn checks_that_could_not_run_are_a_persons_problem() {
1261 let errored = Facts {
1262 check_status: Some(CheckStatus::Errored),
1263 ..facts()
1264 };
1265 assert_eq!(outcome(errored), (Stage::NeedsYou, "wait"));
1266 }
1267
1268 #[test]
1269 fn a_review_asking_for_changes_sends_the_author_back() {
1270 let changes = Facts {
1271 review: reviewed(Some(Verdict::RequestChanges)),
1272 ..facts()
1273 };
1274 assert_eq!(outcome(changes), (Stage::Revising, "revise for review"));
1275 }
1276
1277 #[test]
1278 fn the_author_is_sent_back_only_so_many_times() {
1279 let failing = Facts {
1280 check_status: Some(CheckStatus::Failed),
1281 revisions: MAX_REVISIONS,
1282 ..facts()
1283 };
1284 assert_eq!(outcome(failing), (Stage::NeedsYou, "wait"));
1285 let unconvinced = Facts {
1286 review: reviewed(Some(Verdict::RequestChanges)),
1287 revisions: MAX_REVISIONS,
1288 ..facts()
1289 };
1290 assert_eq!(outcome(unconvinced), (Stage::NeedsYou, "wait"));
1291 // One short of the limit still gets another go.
1292 let once = Facts {
1293 check_status: Some(CheckStatus::Failed),
1294 revisions: MAX_REVISIONS - 1,
1295 ..facts()
1296 };
1297 assert_eq!(outcome(once), (Stage::Revising, "revise for checks"));
1298 }
1299
1300 #[test]
1301 fn a_review_that_could_not_be_written_is_not_retried() {
1302 let broken = Facts {
1303 review: reviewed(None),
1304 ..facts()
1305 };
1306 assert_eq!(outcome(broken), (Stage::NeedsYou, "wait"));
1307 }
1308
1309 #[test]
1310 fn being_behind_only_holds_a_change_up_where_the_repository_says_so() {
1311 let behind = || Facts {
1312 review: reviewed(Some(Verdict::Approve)),
1313 behind: true,
1314 ..facts()
1315 };
1316 // By default it is ready as it is; merging brings it up to date.
1317 assert_eq!(outcome(behind()), (Stage::Ready, "wait"));
1318 let strict = Facts {
1319 require_up_to_date: true,
1320 ..behind()
1321 };
1322 assert_eq!(outcome(strict), (Stage::CatchingUp, "catch up"));
1323 let current = Facts {
1324 review: reviewed(Some(Verdict::Approve)),
1325 ..facts()
1326 };
1327 assert_eq!(outcome(current), (Stage::Ready, "wait"));
1328 }
1329
1330 #[test]
1331 fn a_ready_change_lands_by_itself_only_where_the_repository_says_so() {
1332 let ready = || Facts {
1333 review: reviewed(Some(Verdict::Approve)),
1334 ..facts()
1335 };
1336 assert_eq!(outcome(ready()), (Stage::Ready, "wait"));
1337 let automatic = Facts {
1338 auto_merge: true,
1339 ..ready()
1340 };
1341 assert_eq!(outcome(automatic), (Stage::Ready, "merge"));
1342 // One that is behind is merged too: merging brings it up to date.
1343 let behind = Facts {
1344 auto_merge: true,
1345 behind: true,
1346 ..ready()
1347 };
1348 assert_eq!(outcome(behind), (Stage::Ready, "merge"));
1349 // Unless the repository wants it caught up and checked again first.
1350 let strict = Facts {
1351 auto_merge: true,
1352 behind: true,
1353 require_up_to_date: true,
1354 ..ready()
1355 };
1356 assert_eq!(outcome(strict), (Stage::CatchingUp, "catch up"));
1357 // Nothing short of approved is merged, whatever the setting.
1358 let failing = Facts {
1359 auto_merge: true,
1360 check_status: Some(CheckStatus::Failed),
1361 ..ready()
1362 };
1363 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1364 let unreviewed = Facts {
1365 auto_merge: true,
1366 ..facts()
1367 };
1368 assert_eq!(outcome(unreviewed), (Stage::Reviewing, "review"));
1369 }
1370
1371 #[test]
1372 fn catching_up_reruns_the_checks_but_not_the_review() {
1373 // The merge moved the head, so the checks are waited for again.
1374 let merged_in = Facts {
1375 review: reviewed(Some(Verdict::Approve)),
1376 check_status: None,
1377 ..facts()
1378 };
1379 assert_eq!(outcome(merged_in), (Stage::Checking, "wait"));
1380 }
1381
1382 #[test]
1383 fn a_step_under_way_is_not_started_again() {
1384 for (step, stage) in [
1385 ("review", Stage::Reviewing),
1386 ("revision", Stage::Revising),
1387 ("catch_up", Stage::CatchingUp),
1388 ] {
1389 let busy = Facts {
1390 working_on: Some(step.to_owned()),
1391 // Whatever else is true, the step in hand comes first.
1392 check_status: Some(CheckStatus::Failed),
1393 ..facts()
1394 };
1395 assert_eq!(outcome(busy), (stage, "wait"));
1396 }
1397 let asked = Facts {
1398 review_pending: true,
1399 ..facts()
1400 };
1401 assert_eq!(outcome(asked), (Stage::Reviewing, "wait"));
1402 }
1403
1404 #[test]
1405 fn once_stopped_it_stays_stopped() {
1406 let (lifecycle, next) = decide(Facts {
1407 stalled: Some("The agent could not catch up.".to_owned()),
1408 review: reviewed(Some(Verdict::Approve)),
1409 behind: true,
1410 ..facts()
1411 });
1412 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1413 assert_eq!(lifecycle.detail, "The agent could not catch up.");
1414 assert!(matches!(next, Next::Wait));
1415 }
1416}