Skip to content

g1t/services/work/src/lifecycle.rs

1,921 lines74,337 bytesCodeBlame
1//! Seeing a pull request through. Once a g1t agent has made a change, g1t
2//! takes each remaining step itself: waiting for the checks the
3//! repository's workflows report on it, a review by another agent, sending
4//! the author back to address a failed check (with what its jobs printed)
5//! or the review, and catching up when the branch it would land on has
6//! moved. It stops when the pull request meets everything the default
7//! branch's protection requires, or when it has tried and a person has to
8//! decide.
9//!
10//! This service decides what the next step is and claims it. The runner
11//! service asks, on every event that could change the answer, and carries
12//! the step out in a sandbox.
13
14use g1t_contracts::repos::{GetByIdArgs, Repo, RepoPath};
15use g1t_contracts::time::rfc3339;
16use g1t_contracts::work::*;
17use g1t_contracts::{FailureCode, Membership, Outcome, User, Viewer};
18use g1t_kit::now_ms;
19use std::collections::HashMap;
20
21use serde::Deserialize;
22use worker::Result;
23use worker::wasm_bindgen::JsValue;
24
25use crate::Work;
26use crate::prefetch::Slot;
27use crate::reviews::{AGENT_ID, AGENT_NAME};
28use crate::statuses::{self, WorkflowFacts};
29use crate::rows::ValueRow;
30
31/// How long a claimed step is waited for before it may be taken again.
32const REVIEW_MINUTES: u64 = 20;
33const REVISION_MINUTES: u64 = 60;
34const CATCH_UP_MINUTES: u64 = 30;
35const MERGE_MINUTES: u64 = 2;
36/// Who a merge made by a repository's settings is attributed to. Not an
37/// account: `g1t` cannot be registered.
38pub(crate) const POLICY_ACTOR_ID: &str = "g1t_policy";
39pub(crate) const POLICY_ACTOR_NAME: &str = "g1t";
40/// How much of a failed check's output the author is shown.
41const MAX_CHECK_OUTPUT_CHARS: usize = 4_000;
42/// How many failed jobs' logs the author is shown, and how much of each.
43const MAX_FAILED_JOBS: usize = 3;
44const MAX_JOB_LOG_CHARS: usize = 3_000;
45/// The most pages of a job's log read to find its end.
46const MAX_LOG_PAGES: usize = 6;
47const MANAGED_PAGE: u32 = 200;
48
49/// The part of a pull request's row that tracks its lifecycle.
50#[derive(Deserialize)]
51struct Progress {
52 managed: u8,
53 revisions: u32,
54 revised_at: Option<String>,
55 working_on: Option<String>,
56 working_until: Option<String>,
57 stalled: Option<String>,
58}
59
60/// A merge asked for while the pull request was behind, as stored.
61#[derive(serde::Serialize, Deserialize)]
62#[serde(rename_all = "camelCase")]
63struct LandRequest {
64 actor: User,
65 keep_issue_open: bool,
66}
67
68#[derive(Deserialize)]
69struct LandRow {
70 land_requested: Option<String>,
71 land_requested_at: Option<String>,
72 stalled: Option<String>,
73}
74
75#[derive(Deserialize)]
76struct FinishedReview {
77 finished_at: String,
78 verdict: Option<Verdict>,
79}
80
81#[derive(Deserialize)]
82struct ReviewNote {
83 body: String,
84 path: Option<String>,
85 line: Option<u32>,
86}
87
88/// What the author is being sent back to address.
89pub(crate) enum Feedback {
90 /// The merge queue took it out: its combined state failed.
91 FailedChecks,
92 /// Checks that failed on its head.
93 FailedWorkflows,
94 /// The review that finished at this time.
95 Review(String),
96 /// What a person who asked for changes wrote since the last revision.
97 Person(PersonRequest),
98}
99
100/// A person's request for changes that still stands: their latest verdict,
101/// made after the agent last revised.
102#[derive(Clone, Debug, Deserialize)]
103pub(crate) struct PersonRequest {
104 author_id: String,
105 author_name: String,
106 created_at: String,
107}
108
109/// What should happen next, if it is g1t's turn.
110pub(crate) enum Next {
111 /// A step is under way, or it is a person's turn.
112 Wait,
113 Review,
114 Revise(Feedback),
115 CatchUp,
116 /// Land it, because the repository says ready pull requests land.
117 Merge,
118}
119
120/// Whether g1t made this pull request, and so sees it through.
121pub(crate) fn made_by_g1t(pull: &Pull) -> bool {
122 pull.runtime == Runtime::Hosted && pull.agent == AGENT_NAME && pull.fork.is_some()
123}
124
125/// Whether a verdict by `reviewer` is a person's other than the pull
126/// request's owner (whoever asked g1t for it, or its author): neither
127/// theirs nor g1t's agent's.
128pub(crate) fn from_someone_else(pull: &Pull, reviewer: &str) -> bool {
129 !pull.is_owned_by(reviewer) && reviewer != AGENT_ID
130}
131
132fn at(stage: Stage, detail: impl Into<String>, revisions: u32) -> Lifecycle {
133 Lifecycle {
134 stage,
135 detail: detail.into(),
136 revisions,
137 }
138}
139
140fn times(count: u32) -> String {
141 match count {
142 1 => "once".to_owned(),
143 2 => "twice".to_owned(),
144 count => format!("{count} times"),
145 }
146}
147
148/// Everything the next step depends on.
149struct Facts {
150 /// Still being made: not yet marked ready for review.
151 draft: bool,
152 /// Why g1t stopped, if it has.
153 stalled: Option<String>,
154 /// The step under way, if one was claimed and is still being waited for.
155 working_on: Option<String>,
156 /// `failed` when the merge queue took it out.
157 check_status: Option<CheckStatus>,
158 /// A review someone asked for is being written.
159 review_pending: bool,
160 revisions: u32,
161 /// The latest finished review of the change as it is now.
162 review: Option<FinishedReview>,
163 /// Whether the branch it would land on has moved without it.
164 behind: bool,
165 /// Whether it is known to conflict with the branch it would land on.
166 conflicting: bool,
167 /// Whether the repository lands a ready pull request by itself.
168 auto_merge: bool,
169 /// Whether the repository refuses to merge one that is behind.
170 require_up_to_date: bool,
171 /// Whether a second agent reviews it without being asked.
172 agent_review: bool,
173 /// How many times the author may be sent back.
174 max_revisions: u32,
175 /// What the repository's approval rule still wants, if anything.
176 approvals_missing: Option<String>,
177 /// A person asked for changes since the agent last revised.
178 person_request: Option<PersonRequest>,
179 /// Its place in the merge queue, and what is ahead of it there.
180 queued: Option<(QueueState, Vec<u32>)>,
181 /// What the checks on its head say, against the required ones.
182 workflows: WorkflowFacts,
183 /// Why the agent's change has low confidence, when the repository asks
184 /// a person before merging one and no person has approved it since.
185 low_confidence: Option<String>,
186}
187
188/// Where a pull request stands, and the step to take if it is g1t's turn.
189///
190/// The order is: nothing while a step is under way; a person asking for
191/// changes is answered first; its checks must finish, and the required
192/// ones pass; then a review must approve; then it must be up to date. A
193/// person's or an agent's request for changes, or a failed check, sends the
194/// author back, a limited number of times, after which a person is asked.
195/// A check the branch does not require stops holding it once the
196/// revisions run out.
197fn decide(facts: Facts) -> (Lifecycle, Next) {
198 let revisions = facts.revisions;
199 let wait = |stage, detail: &str| (at(stage, detail, revisions), Next::Wait);
200 let exhausted = revisions >= facts.max_revisions;
201
202 if facts.draft {
203 return wait(Stage::Working, "g1t is making the change.");
204 }
205 if let Some(reason) = &facts.stalled {
206 return wait(Stage::NeedsYou, reason);
207 }
208 if let Some((state, ahead)) = &facts.queued {
209 let named = ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ");
210 let detail = match (state, ahead.is_empty()) {
211 (QueueState::Testing, true) => "In the merge queue: being tested on the default branch as it is.".to_owned(),
212 (QueueState::Testing, false) => {
213 format!("In the merge queue: being tested together with {named}, ahead of it.")
214 }
215 (QueueState::Passed, true) => "Passed in the merge queue. Landing.".to_owned(),
216 (QueueState::Passed, false) => {
217 format!("Passed in the merge queue together with {named}. It lands once they have.")
218 }
219 _ => "In the merge queue, waiting for its turn to be tested.".to_owned(),
220 };
221 return wait(Stage::Queued, &detail);
222 }
223 match facts.working_on.as_deref() {
224 Some("revision") => {
225 return wait(
226 Stage::Revising,
227 "The agent is addressing what the checks or the review found.",
228 );
229 }
230 Some("catch_up") => {
231 return wait(
232 Stage::CatchingUp,
233 "The agent is merging in the branch this will land on, which has moved.",
234 );
235 }
236 Some("answer") => {
237 return wait(
238 Stage::Answering,
239 "The agent is answering what another agent asked it.",
240 );
241 }
242 Some("merge") => return wait(Stage::Ready, "Merging."),
243 Some(_) => return wait(Stage::Reviewing, "g1t is reviewing the change."),
244 None => {}
245 }
246 if facts.review_pending {
247 return wait(Stage::Reviewing, "g1t is reviewing the change.");
248 }
249 // A person asked for changes: the agent makes them, as it would for a
250 // review it asked for, before anything else.
251 if let Some(request) = &facts.person_request {
252 if exhausted {
253 return wait(
254 Stage::NeedsYou,
255 &format!(
256 "{} asked for changes, and the agent has already revised {}.",
257 request.author_name,
258 times(revisions)
259 ),
260 );
261 }
262 return (
263 at(
264 Stage::Revising,
265 format!(
266 "{} asked for changes. The agent is being sent back to make them.",
267 request.author_name
268 ),
269 revisions,
270 ),
271 Next::Revise(Feedback::Person(request.clone())),
272 );
273 }
274
275 // The merge queue took it out: its change failed together with what
276 // was ahead of it.
277 match facts.check_status {
278 Some(CheckStatus::Failed) if exhausted => {
279 return wait(
280 Stage::NeedsYou,
281 &format!("It failed in the merge queue after the agent revised {}.", times(revisions)),
282 );
283 }
284 Some(CheckStatus::Failed) => {
285 return (
286 at(
287 Stage::Revising,
288 "It failed in the merge queue. The agent is being sent back to fix it.",
289 revisions,
290 ),
291 Next::Revise(Feedback::FailedChecks),
292 );
293 }
294 Some(CheckStatus::Errored) => {
295 return wait(Stage::NeedsYou, "Its checks could not be run.");
296 }
297 _ => {}
298 }
299
300 // Its checks: the agent fixes any that failed. Once it is out of
301 // revisions, only the checks the branch requires still hold it.
302 if !facts.workflows.failed.is_empty() {
303 let failed = statuses::list(&facts.workflows.failed);
304 let required = facts.workflows.required_failed();
305 if !exhausted {
306 return (
307 at(
308 Stage::Revising,
309 format!("{failed} failed. The agent is being sent back to fix it."),
310 revisions,
311 ),
312 Next::Revise(Feedback::FailedWorkflows),
313 );
314 }
315 if !required.is_empty() {
316 return wait(
317 Stage::NeedsYou,
318 &format!(
319 "The required {} {} still {} after the agent revised {}.",
320 if required.len() == 1 { "check" } else { "checks" },
321 statuses::list(&required),
322 if required.len() == 1 { "fails" } else { "fail" },
323 times(revisions)
324 ),
325 );
326 }
327 }
328 if !facts.workflows.pending.is_empty() {
329 return wait(
330 Stage::Checking,
331 &format!("Waiting for {} to finish.", statuses::list(&facts.workflows.pending)),
332 );
333 }
334 let expected = facts.workflows.expected();
335 if !expected.is_empty() {
336 return wait(
337 Stage::Checking,
338 &format!(
339 "Waiting for the required {} {} to report on its latest commit.",
340 if expected.len() == 1 { "check" } else { "checks" },
341 statuses::list(&expected)
342 ),
343 );
344 }
345
346 // A second agent reviews it, unless the repository leaves review to people.
347 if facts.agent_review {
348 match facts.review {
349 None => {
350 return (
351 at(
352 Stage::Reviewing,
353 "g1t is about to review the change.",
354 revisions,
355 ),
356 Next::Review,
357 );
358 }
359 Some(FinishedReview { verdict: None, .. }) => {
360 return wait(Stage::NeedsYou, "The review could not be completed.");
361 }
362 Some(FinishedReview {
363 verdict: Some(Verdict::RequestChanges),
364 ..
365 }) if exhausted => {
366 return wait(
367 Stage::NeedsYou,
368 &format!(
369 "The review still asks for changes after the agent revised {}.",
370 times(revisions)
371 ),
372 );
373 }
374 Some(FinishedReview {
375 verdict: Some(Verdict::RequestChanges),
376 finished_at,
377 }) => {
378 return (
379 at(
380 Stage::Revising,
381 "The review asked for changes. The agent is being sent back to make them.",
382 revisions,
383 ),
384 Next::Revise(Feedback::Review(finished_at)),
385 );
386 }
387 Some(FinishedReview {
388 verdict: Some(Verdict::Approve),
389 ..
390 }) => {}
391 }
392 }
393
394 // A conflict found ahead of time is resolved before anything else that
395 // is left: it could not merge, by a person or by the queue, until then.
396 if facts.conflicting {
397 return (
398 at(
399 Stage::CatchingUp,
400 "It conflicts with the branch it will land on. The agent is merging that branch in and resolving the conflicts.",
401 revisions,
402 ),
403 Next::CatchUp,
404 );
405 }
406 // Only where the repository insists is catching up a step of its own,
407 // followed by its checks again. Elsewhere it happens as part of merging.
408 if facts.behind && facts.require_up_to_date {
409 return (
410 at(
411 Stage::CatchingUp,
412 "The branch it will land on has moved. The agent is catching up.",
413 revisions,
414 ),
415 Next::CatchUp,
416 );
417 }
418 // The repository wants approvals this does not have yet: people's turn,
419 // so it is shown as needing someone, not as g1t still working.
420 if let Some(missing) = &facts.approvals_missing {
421 return wait(Stage::NeedsYou, missing);
422 }
423 // Everything else is met, but g1t is not sure of the change: a person
424 // decides, rather than auto-merge or the queue.
425 if let Some(reasons) = &facts.low_confidence {
426 return wait(
427 Stage::NeedsYou,
428 &format!(
429 "The agent's confidence in this change is low ({reasons}). This repository asks a person before merging it: approve it to let it land, or ask for changes."
430 ),
431 );
432 }
433 if facts.auto_merge {
434 return (
435 at(
436 Stage::Ready,
437 "Everything this repository asks for is met. Merging, as its settings say.",
438 revisions,
439 ),
440 Next::Merge,
441 );
442 }
443 wait(
444 Stage::Ready,
445 if facts.behind {
446 "Ready to merge. Merging brings it up to date with the default branch first."
447 } else {
448 "Everything this repository asks for is met. Ready to merge."
449 },
450 )
451}
452
453impl Work {
454 /// Where a pull request stands and what g1t does next, remembered so
455 /// lists can show it without working it out again. `None` for one g1t
456 /// is not seeing through.
457 pub(crate) async fn assess(
458 &self,
459 pull: &Pull,
460 issue: &Option<Issue>,
461 behind: bool,
462 ) -> Result<Option<(Lifecycle, Next)>> {
463 Ok(self
464 .assess_with_confidence(pull, issue, behind)
465 .await?
466 .map(|(lifecycle, next, _)| (lifecycle, next)))
467 }
468
469 /// [`Self::assess`], with how sure g1t is of the change once the agent
470 /// has finished it.
471 pub(crate) async fn assess_with_confidence(
472 &self,
473 pull: &Pull,
474 issue: &Option<Issue>,
475 behind: bool,
476 ) -> Result<Option<(Lifecycle, Next, Option<Confidence>)>> {
477 let assessed = self.assess_now(pull, issue, behind).await?;
478 if let Some((lifecycle, _, _)) = &assessed
479 && !self.remembered_as(&pull.id, lifecycle)?
480 {
481 self.remember(&pull.id, lifecycle).await?;
482 }
483 Ok(assessed)
484 }
485
486 /// Whether the row read for this request already says `lifecycle`, so
487 /// that showing a pull request does not write it again unchanged.
488 fn remembered_as(&self, pull_id: &str, lifecycle: &Lifecycle) -> Result<bool> {
489 let Some(found) = self.prefetched_pull(pull_id) else {
490 return Ok(false);
491 };
492 Ok(found
493 .first::<crate::rows::Snapshot>(Slot::Pull)?
494 .is_some_and(|stored| {
495 stored.stage == Some(lifecycle.stage) && stored.stage_detail.as_deref() == Some(lifecycle.detail.as_str())
496 }))
497 }
498
499 /// Saves where a pull request stands, for [`Self::remembered`].
500 pub(crate) async fn remember(&self, pull_id: &str, lifecycle: &Lifecycle) -> Result<()> {
501 let stage = serde_json::to_value(lifecycle.stage)?;
502 self.db
503 .prepare("UPDATE pulls SET stage = ?, stage_detail = ? WHERE id = ?")
504 .bind(&[
505 stage.as_str().unwrap_or_default().into(),
506 lifecycle.detail.as_str().into(),
507 pull_id.into(),
508 ])?
509 .run()
510 .await?;
511 Ok(())
512 }
513
514 async fn assess_now(
515 &self,
516 pull: &Pull,
517 // What done means is in its body, for the agent; the merge waits on
518 // the branch's required checks, not on the issue.
519 _issue: &Option<Issue>,
520 behind: bool,
521 ) -> Result<Option<(Lifecycle, Next, Option<Confidence>)>> {
522 if !pull.status.is_active() {
523 return Ok(None);
524 }
525 // As read for this request (prefetch.rs), or now.
526 let prefetched = self.prefetched_pull(&pull.id);
527 let progress = match &prefetched {
528 Some(found) => found.first::<Progress>(Slot::Pull)?,
529 None => {
530 self.db
531 .prepare(
532 "SELECT managed, revisions, revised_at, working_on, working_until, stalled
533 FROM pulls WHERE id = ?",
534 )
535 .bind(&[pull.id.as_str().into()])?
536 .first::<Progress>(None)
537 .await?
538 }
539 };
540 let Some(progress) = progress.filter(|progress| progress.managed != 0) else {
541 return Ok(None);
542 };
543 let now = rfc3339(now_ms());
544 let working_on = progress
545 .working_until
546 .as_deref()
547 .is_some_and(|until| until > now.as_str())
548 .then(|| progress.working_on.clone().unwrap_or_default());
549 let review = match &prefetched {
550 Some(found) => found.first::<FinishedReview>(Slot::Review)?,
551 None => {
552 self.db
553 .prepare(
554 "SELECT finished_at, verdict FROM review_runs
555 WHERE pull_id = ? AND finished_at IS NOT NULL ORDER BY id DESC LIMIT 1",
556 )
557 .bind(&[pull.id.as_str().into()])?
558 .first::<FinishedReview>(None)
559 .await?
560 }
561 }
562 // A review of what the change was before its last revision says
563 // nothing about what it is now.
564 .filter(|review| {
565 progress
566 .revised_at
567 .as_deref()
568 .is_none_or(|revised| review.finished_at.as_str() >= revised)
569 });
570 let settings = self.settings_for(pull).await?;
571 let workflows = WorkflowFacts::of(
572 &self.statuses(&pull.repo_id, pull.head_commit.as_deref()).await?,
573 &settings.required_checks,
574 );
575 // Once the agent has finished the change: how sure g1t is of it, and
576 // whether that holds it for a person.
577 let (confidence, low_confidence) = if pull.status == PullStatus::Draft {
578 (None, None)
579 } else {
580 let review_comments = match &review {
581 Some(review) => self.review_comments(pull, &review.finished_at).await?,
582 None => 0,
583 };
584 let confidence = self
585 .assess_confidence(
586 pull,
587 crate::confidence::Signals {
588 required: crate::confidence::RequiredSignal::of(&workflows.required),
589 queue_failed: pull.check_status == Some(CheckStatus::Failed),
590 revisions: progress.revisions,
591 agent_review: settings.agent_review,
592 review: review.as_ref().and_then(|review| review.verdict),
593 review_comments,
594 ..Default::default()
595 },
596 )
597 .await?;
598 let held = settings.hold_low_confidence
599 && confidence.level == ConfidenceLevel::Low
600 && !self.person_approved(pull, progress.revised_at.as_deref()).await?;
601 let reasons = held.then(|| confidence.reasons.join(", "));
602 (Some(confidence), reasons)
603 };
604 let (lifecycle, next) = decide(Facts {
605 draft: pull.status == PullStatus::Draft,
606 stalled: progress.stalled,
607 working_on,
608 check_status: pull.check_status,
609 review_pending: self.review_pending(&pull.id).await?,
610 revisions: progress.revisions,
611 review,
612 behind,
613 conflicting: self.conflicting_files(pull).await?.is_some(),
614 auto_merge: settings.auto_merge,
615 require_up_to_date: settings.require_up_to_date,
616 agent_review: settings.agent_review,
617 max_revisions: settings.max_revisions,
618 approvals_missing: self.approvals_gap(&settings, pull).await?,
619 person_request: self
620 .person_request(pull, progress.revised_at.as_deref())
621 .await?,
622 queued: self.queued_entry(&pull.id).await?,
623 workflows,
624 low_confidence,
625 });
626 Ok(Some((lifecycle, next, confidence)))
627 }
628
629 /// The latest request for changes by a person other than its owner
630 /// (whoever asked g1t for it, or its author), if it is that person's
631 /// latest verdict and came after the last revision.
632 async fn person_request(
633 &self,
634 pull: &Pull,
635 revised_at: Option<&str>,
636 ) -> Result<Option<PersonRequest>> {
637 #[derive(Deserialize)]
638 struct Verdicts {
639 author_id: String,
640 author_name: String,
641 verdict: Verdict,
642 created_at: String,
643 }
644 let rows = match self.prefetched_pull(&pull.id) {
645 Some(found) => found
646 .rows::<Verdicts>(Slot::Verdicts)?
647 .into_iter()
648 .filter(|row| from_someone_else(pull, &row.author_id))
649 .collect(),
650 None => self
651 .db
652 .prepare(
653 "SELECT author_id, author_name, verdict, created_at FROM comments
654 WHERE repo_id = ? AND number = ? AND verdict IS NOT NULL
655 AND author_id != ? AND author_id != ?
656 ORDER BY id",
657 )
658 .bind(&[
659 pull.repo_id.as_str().into(),
660 pull.number.into(),
661 pull.owner().id.as_str().into(),
662 AGENT_ID.into(),
663 ])?
664 .all()
665 .await?
666 .results::<Verdicts>()?,
667 };
668 // Each person's latest verdict is the one that stands.
669 let mut latest: HashMap<String, Verdicts> = HashMap::new();
670 for row in rows {
671 latest.insert(row.author_id.clone(), row);
672 }
673 Ok(latest
674 .into_values()
675 .filter(|row| row.verdict == Verdict::RequestChanges)
676 .filter(|row| revised_at.is_none_or(|revised| row.created_at.as_str() > revised))
677 .max_by(|a, b| a.created_at.cmp(&b.created_at))
678 .map(|row| PersonRequest {
679 author_id: row.author_id,
680 author_name: row.author_name,
681 created_at: row.created_at,
682 }))
683 }
684
685 /// Marks a pull request a g1t agent has just opened as one g1t sees
686 /// through.
687 pub(crate) async fn manage(&self, pull: &Pull) -> Result<()> {
688 if !made_by_g1t(pull) {
689 return Ok(());
690 }
691 self.db
692 .prepare("UPDATE pulls SET managed = 1 WHERE id = ?")
693 .bind(&[pull.id.as_str().into()])?
694 .run()
695 .await?;
696 Ok(())
697 }
698
699 /// Takes a step for a pull request, if nobody else has. One statement,
700 /// so that two callers cannot both take it.
701 pub(crate) async fn claim(&self, pull_id: &str, step: &str, minutes: u64, revising: bool) -> Result<bool> {
702 let now = now_ms();
703 let revision = if revising {
704 ", revisions = revisions + 1, revised_at = ?1"
705 } else {
706 ""
707 };
708 Ok(self
709 .db
710 .prepare(format!(
711 "UPDATE pulls SET working_on = ?2, working_until = ?3{revision}
712 WHERE id = ?4 AND status = 'open' AND stalled IS NULL
713 AND (working_until IS NULL OR working_until < ?1)
714 RETURNING id AS value"
715 ))
716 .bind(&[
717 rfc3339(now).into(),
718 step.into(),
719 rfc3339(now + minutes * 60 * 1000).into(),
720 pull_id.into(),
721 ])?
722 .first::<ValueRow>(None)
723 .await?
724 .is_some())
725 }
726
727 /// What the author is told when sent back: the checks that failed and
728 /// what their failing jobs printed, why the merge queue took it out, or
729 /// the review and its comments on lines.
730 async fn feedback(&self, pull: &Pull, feedback: &Feedback) -> Result<String> {
731 match feedback {
732 Feedback::FailedChecks => {
733 let run = self.latest_checks(&pull.id).await?;
734 // Why it failed, when that is more than a list of commands:
735 // the merge queue saying what broke in the combined state.
736 let why = run.as_ref().and_then(|run| run.error.clone()).map(|error| format!("{error}\n\n")).unwrap_or_default();
737 let failed: Vec<String> = run
738 .map(|run| run.results)
739 .unwrap_or_default()
740 .into_iter()
741 .filter(|result| !result.passed)
742 .map(|result| {
743 let length = result.output.chars().count();
744 let output: String = result
745 .output
746 .chars()
747 .skip(length.saturating_sub(MAX_CHECK_OUTPUT_CHARS))
748 .collect();
749 let exit = result
750 .exit_code
751 .map_or("it was stopped for taking too long".to_owned(), |code| {
752 format!("exit code {code}")
753 });
754 format!("`{}` failed ({exit}):\n\n{}", result.command, output.trim())
755 })
756 .collect();
757 if failed.is_empty() {
758 return Ok(format!(
759 "{why}Find the cause, fix it in your change, and push. Use get_workflow_run and get_job_logs for any workflow named above."
760 ));
761 }
762 Ok(format!(
763 "{why}These commands failed against your change.\n\n{}",
764 failed.join("\n\n")
765 ))
766 }
767 Feedback::FailedWorkflows => {
768 let (statuses, settings) = futures_util::future::try_join(
769 self.statuses(&pull.repo_id, pull.head_commit.as_deref()),
770 self.settings_for(pull),
771 )
772 .await?;
773 let required = |context: &str| {
774 let name = g1t_contracts::work::check_name(context).0;
775 settings.required_checks.iter().any(|wanted| wanted.eq_ignore_ascii_case(name))
776 };
777 let failing: Vec<&CommitStatus> =
778 statuses.iter().filter(|s| s.state == "failure" || s.state == "error").collect();
779 let failed: Vec<String> = failing
780 .iter()
781 .map(|s| {
782 let run = s.target_url.as_deref().and_then(|url| url.rsplit('/').next()).unwrap_or_default();
783 format!(
784 "- {}{} ({}): run `{run}`, {}",
785 s.context,
786 if required(&s.context) { ", required to merge" } else { "" },
787 s.description.as_deref().unwrap_or("failed"),
788 s.target_url.as_deref().unwrap_or_default()
789 )
790 })
791 .collect();
792 let logs = self.failing_logs(pull, &failing).await.unwrap_or_default();
793 // Named outright: the agent cannot guess it from its fork.
794 let repo = g1t_kit::call::<_, Option<RepoPath>>(
795 &self.repos,
796 "path_by_id",
797 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
798 )
799 .await?
800 .map(|path| format!("{}/{}", path.namespace, path.name))
801 .unwrap_or_default();
802 let logs = if logs.is_empty() {
803 String::new()
804 } else {
805 format!("\n\nThe end of what the failing jobs printed:\n\n{}", logs.join("\n\n"))
806 };
807 Ok(format!(
808 "These checks failed on your latest commit to {repo}. They are the repository's workflows, run on your pull request:\n\n{}{logs}\n\n\
809 For more, use the `get_workflow_run` tool (repo `{repo}` and the run's id), \
810 then `get_job_logs` for the job that failed. Fix the cause in the code, not the workflow, \
811 unless the workflow itself is wrong. Push, and the checks run again.",
812 failed.join("\n")
813 ))
814 }
815 Feedback::Review(finished_at) => {
816 // Everything a review says is recorded at the moment it finished.
817 let notes = self
818 .db
819 .prepare(
820 "SELECT body, path, line FROM comments
821 WHERE repo_id = ? AND number = ? AND author_id = ? AND created_at = ?
822 ORDER BY id",
823 )
824 .bind(&[
825 pull.repo_id.as_str().into(),
826 pull.number.into(),
827 AGENT_ID.into(),
828 finished_at.as_str().into(),
829 ])?
830 .all()
831 .await?
832 .results::<ReviewNote>()?;
833 let mut on_lines = Vec::new();
834 let mut summary = String::new();
835 for note in notes {
836 match (note.path, note.line) {
837 (Some(path), Some(line)) => {
838 on_lines.push(format!("- `{path}` line {line}: {}", note.body));
839 }
840 (Some(path), None) => on_lines.push(format!("- `{path}`: {}", note.body)),
841 (None, _) => summary = note.body,
842 }
843 }
844 let mut text = format!(
845 "Another agent reviewed your change and asked for changes.\n\n{summary}"
846 );
847 if !on_lines.is_empty() {
848 text.push_str("\n\nIts comments on lines:\n");
849 text.push_str(&on_lines.join("\n"));
850 }
851 Ok(text)
852 }
853 Feedback::Person(request) => {
854 // What they wrote since the agent last revised, which their
855 // request for changes closes.
856 let revised: Option<String> = self
857 .db
858 .prepare("SELECT revised_at AS value FROM pulls WHERE id = ?")
859 .bind(&[pull.id.as_str().into()])?
860 .first::<Option<String>>(Some("value"))
861 .await?
862 .flatten();
863 let notes = self
864 .db
865 .prepare(
866 "SELECT body, path, line FROM comments
867 WHERE repo_id = ? AND number = ? AND author_id = ?
868 AND created_at > ? AND created_at <= ?
869 ORDER BY id",
870 )
871 .bind(&[
872 pull.repo_id.as_str().into(),
873 pull.number.into(),
874 request.author_id.as_str().into(),
875 revised.unwrap_or_default().into(),
876 request.created_at.as_str().into(),
877 ])?
878 .all()
879 .await?
880 .results::<ReviewNote>()?;
881 let mut on_lines = Vec::new();
882 let mut said = Vec::new();
883 for note in notes {
884 match (note.path, note.line) {
885 (Some(path), Some(line)) => {
886 on_lines.push(format!("- `{path}` line {line}: {}", note.body));
887 }
888 (Some(path), None) => on_lines.push(format!("- `{path}`: {}", note.body)),
889 (None, _) => said.push(note.body),
890 }
891 }
892 let mut text = format!(
893 "{} reviewed your change and asked for changes.\n\n{}",
894 request.author_name,
895 said.join("\n\n")
896 );
897 if !on_lines.is_empty() {
898 text.push_str("\n\nTheir comments on lines:\n");
899 text.push_str(&on_lines.join("\n"));
900 }
901 Ok(text)
902 }
903 }
904 }
905
906 /// The end of what the failed jobs of failing workflow runs printed, a
907 /// few jobs at most, for an agent sent back to fix them. Empty where the
908 /// runs or their logs cannot be read.
909 async fn failing_logs(&self, pull: &Pull, failing: &[&CommitStatus]) -> Result<Vec<String>> {
910 use g1t_contracts::actions::{JobLog, LogsArgs, RunArgs, RunDetail};
911 let Some(repo) = g1t_kit::call::<_, Option<RepoPath>>(
912 &self.repos,
913 "path_by_id",
914 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
915 )
916 .await?
917 else {
918 return Ok(Vec::new());
919 };
920 let viewer = self.owner_viewer(pull).await?;
921 let mut out = Vec::new();
922 for status in failing {
923 let Some(run_id) = status
924 .target_url
925 .as_deref()
926 .and_then(|url| url.split("/actions/runs/").nth(1))
927 .map(|rest| rest.split(['/', '?', '#']).next().unwrap_or_default().to_owned())
928 .filter(|id| !id.is_empty())
929 else {
930 continue;
931 };
932 let detail: Outcome<RunDetail> = g1t_kit::call(
933 &self.actions,
934 "run",
935 &RunArgs { repo: repo.clone(), viewer: viewer.clone(), id: run_id },
936 )
937 .await?;
938 let Outcome::Ok(detail) = detail else { continue };
939 let failed_jobs = detail.jobs.into_iter().filter(|job| {
940 job.status == "completed" && matches!(job.conclusion.as_deref(), Some("failure" | "timed_out"))
941 });
942 for job in failed_jobs {
943 if out.len() >= MAX_FAILED_JOBS {
944 return Ok(out);
945 }
946 let mut text = String::new();
947 let mut after = 0;
948 for _ in 0..MAX_LOG_PAGES {
949 let page: Outcome<JobLog> = g1t_kit::call(
950 &self.actions,
951 "logs",
952 &LogsArgs { repo: repo.clone(), viewer: viewer.clone(), job: job.id.clone(), after },
953 )
954 .await?;
955 let Outcome::Ok(page) = page else { break };
956 let Some(last) = page.chunks.last().map(|chunk| chunk.seq) else { break };
957 for chunk in &page.chunks {
958 text.push_str(&chunk.text);
959 if !chunk.text.ends_with('\n') {
960 text.push('\n');
961 }
962 }
963 // Only the end is kept, so the start can go as it is read.
964 let length = text.chars().count();
965 if length > MAX_JOB_LOG_CHARS * 2 {
966 text = text.chars().skip(length - MAX_JOB_LOG_CHARS).collect();
967 }
968 after = last;
969 if page.chunks.len() < 500 {
970 break;
971 }
972 }
973 let length = text.chars().count();
974 let tail: String = text.chars().skip(length.saturating_sub(MAX_JOB_LOG_CHARS)).collect();
975 if !tail.trim().is_empty() {
976 out.push(format!("{} / {}:\n```\n{}\n```", status.context, job.name, tail.trim_end()));
977 }
978 }
979 }
980 Ok(out)
981 }
982
983 pub(crate) async fn advance(&self, a: AdvanceArgs) -> Result<Advance> {
984 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
985 return Ok(Advance::None);
986 };
987 if pull.status != PullStatus::Open {
988 return Ok(Advance::None);
989 }
990 // As a member: a private repository would look missing otherwise,
991 // and the pull request would never move.
992 let viewer: Viewer = self.owner_viewer(&pull).await?;
993 let repo: Outcome<Repo> = g1t_kit::call(
994 &self.repos,
995 "get_by_id",
996 &GetByIdArgs {
997 id: pull.repo_id.clone(),
998 viewer,
999 },
1000 )
1001 .await?;
1002 let (Outcome::Ok(repo), Some(source)) = (crate::retired::unless_archived(repo), pull.fork.clone()) else {
1003 return Ok(Advance::None);
1004 };
1005 // The branch it merges into: what it catches up with.
1006 let base = pull.base_branch(&repo.default_branch).to_owned();
1007 let issue = match pull.issue {
1008 Some(number) => self.issue(&pull.repo_id, number).await?,
1009 None => None,
1010 };
1011 let behind = self.is_behind(&repo.id, &pull).await?;
1012 let Some((lifecycle, next)) = self.assess(&pull, &issue, behind).await? else {
1013 return Ok(Advance::None);
1014 };
1015
1016 if matches!(next, Next::Merge) {
1017 self.merge_by_policy(&repo, &pull).await?;
1018 return Ok(Advance::None);
1019 }
1020 let (step, minutes) = match &next {
1021 Next::Wait | Next::Merge => return Ok(Advance::None),
1022 Next::Review => ("review", REVIEW_MINUTES),
1023 Next::Revise(_) => ("revision", REVISION_MINUTES),
1024 Next::CatchUp => ("catch_up", CATCH_UP_MINUTES),
1025 };
1026 let feedback = match &next {
1027 Next::Revise(feedback) => self.feedback(&pull, feedback).await?,
1028 Next::CatchUp => self.conflict_note(&pull, &base).await?,
1029 _ => String::new(),
1030 };
1031 if !self
1032 .claim(&pull.id, step, minutes, matches!(next, Next::Revise(_)))
1033 .await?
1034 {
1035 return Ok(Advance::None);
1036 }
1037 // Said in the conversation, so nobody has to wonder why a review or
1038 // a new commit appeared.
1039 let told = match &next {
1040 Next::Review => {
1041 self.db
1042 .prepare(
1043 "UPDATE pulls SET reviewers = json_insert(reviewers, '$[#]', ?1)
1044 WHERE id = ?2 AND NOT EXISTS (
1045 SELECT 1 FROM json_each(pulls.reviewers) WHERE json_each.value = ?1)",
1046 )
1047 .bind(&[AGENT_NAME.into(), pull.id.as_str().into()])?
1048 .run()
1049 .await?;
1050 "requested a review from g1t".to_owned()
1051 }
1052 Next::Revise(Feedback::FailedChecks) => {
1053 "sent g1t back to fix what failed in the merge queue".to_owned()
1054 }
1055 Next::Revise(Feedback::FailedWorkflows) => {
1056 "sent g1t back to fix the failed checks".to_owned()
1057 }
1058 Next::Revise(Feedback::Review(_)) => {
1059 "sent g1t back to address the review".to_owned()
1060 }
1061 _ => format!("asked g1t to bring this up to date with {base}"),
1062 };
1063 self.note(
1064 &pull.repo_id,
1065 pull.number,
1066 (POLICY_ACTOR_ID, POLICY_ACTOR_NAME),
1067 &told,
1068 )
1069 .await?;
1070 let job = LifecycleJob {
1071 pull_id: pull.id,
1072 repo: RepoPath {
1073 namespace: repo.namespace,
1074 name: repo.name,
1075 },
1076 number: pull.number,
1077 author: pull.requested_by.unwrap_or(pull.author),
1078 source,
1079 branch: None,
1080 default_branch: base,
1081 title: pull.title,
1082 description: pull.body.unwrap_or_default(),
1083 issue,
1084 feedback,
1085 round: lifecycle.revisions + 1,
1086 };
1087 Ok(match next {
1088 Next::Review => Advance::Review { job },
1089 Next::Revise(_) => Advance::Revise { job },
1090 Next::CatchUp => Advance::CatchUp { job },
1091 Next::Wait | Next::Merge => Advance::None,
1092 })
1093 }
1094
1095 /// Lands a pull request that is ready, on the authority of the
1096 /// repository's settings instead of a person's click.
1097 async fn merge_by_policy(&self, repo: &Repo, pull: &Pull) -> Result<()> {
1098 if !self.claim(&pull.id, "merge", MERGE_MINUTES, false).await? {
1099 return Ok(());
1100 }
1101 // g1t acts for the workspace whose members turned this on, with the
1102 // default base permission (Write), which merging needs.
1103 let actor = User {
1104 id: POLICY_ACTOR_ID.to_owned(),
1105 username: POLICY_ACTOR_NAME.to_owned(),
1106 verified: true,
1107 workspaces: vec![Membership::member(repo.namespace.to_lowercase())],
1108 ..User::default()
1109 };
1110 let merged = self
1111 .merge_pull(PullActionArgs {
1112 actor,
1113 repo: RepoPath {
1114 namespace: repo.namespace.clone(),
1115 name: repo.name.clone(),
1116 },
1117 number: pull.number,
1118 summary: String::new(),
1119 keep_issue_open: false,
1120 ignore_checks: false,
1121 })
1122 .await?;
1123 match merged {
1124 Outcome::Ok(_) => Ok(()),
1125 // Most likely the branch moved in the moment between: let go, and
1126 // the next look at it will catch up and try again.
1127 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
1128 self.db
1129 .prepare(
1130 "UPDATE pulls SET working_on = NULL, working_until = NULL
1131 WHERE id = ? AND working_on = 'merge'",
1132 )
1133 .bind(&[pull.id.as_str().into()])?
1134 .run()
1135 .await?;
1136 Ok(())
1137 }
1138 Outcome::Fail(failure) => {
1139 self.stall(StallArgs {
1140 pull_id: pull.id.clone(),
1141 reason: format!("g1t could not merge this: {}", failure.message),
1142 by: None,
1143 })
1144 .await?;
1145 Ok(())
1146 }
1147 }
1148 }
1149
1150 /// Records that a merge was asked for while the pull request was
1151 /// behind, and announces it so that the runner brings it up to date.
1152 pub(crate) async fn request_landing(
1153 &self,
1154 pull: &Pull,
1155 actor: &User,
1156 keep_issue_open: bool,
1157 ) -> Result<()> {
1158 let now = now_ms();
1159 let request = serde_json::to_string(&LandRequest {
1160 actor: actor.clone(),
1161 keep_issue_open,
1162 })?;
1163 self.db
1164 .prepare(
1165 "UPDATE pulls
1166 SET land_requested = ?, land_requested_at = ?, stalled = NULL,
1167 working_on = 'catch_up', working_until = ?
1168 WHERE id = ?",
1169 )
1170 .bind(&[
1171 request.into(),
1172 rfc3339(now).into(),
1173 rfc3339(now + CATCH_UP_MINUTES * 60 * 1000).into(),
1174 pull.id.as_str().into(),
1175 ])?
1176 .run()
1177 .await?;
1178 self.publish(
1179 "pull.merge_requested",
1180 &pull.repo_id,
1181 actor,
1182 Self::pull_event(pull),
1183 )
1184 .await
1185 }
1186
1187 /// The merge waiting on a pull request, if one was asked for recently
1188 /// enough to still stand.
1189 async fn land_request(&self, pull_id: &str) -> Result<Option<LandRequest>> {
1190 let row = self.land_row(pull_id).await?;
1191 let oldest = rfc3339(now_ms().saturating_sub(CATCH_UP_MINUTES * 60 * 1000));
1192 Ok(row
1193 .filter(|row| {
1194 row.land_requested_at
1195 .as_deref()
1196 .is_some_and(|at| at >= oldest.as_str())
1197 })
1198 .and_then(|row| row.land_requested)
1199 .and_then(|request| serde_json::from_str(&request).ok()))
1200 }
1201
1202 /// Whether a merge is waiting on a pull request, and why g1t stopped
1203 /// working on it if it did.
1204 pub(crate) async fn landing_state(&self, pull_id: &str) -> Result<(bool, Option<String>)> {
1205 let stalled = self.land_row(pull_id).await?.and_then(|row| row.stalled);
1206 Ok((self.land_request(pull_id).await?.is_some(), stalled))
1207 }
1208
1209 /// The merge waiting on a pull request and why g1t stopped, as read
1210 /// for this request (prefetch.rs) or now.
1211 async fn land_row(&self, pull_id: &str) -> Result<Option<LandRow>> {
1212 if let Some(found) = self.prefetched_pull(pull_id) {
1213 return found.first::<LandRow>(Slot::Pull);
1214 }
1215 self.db
1216 .prepare("SELECT land_requested, land_requested_at, stalled FROM pulls WHERE id = ?")
1217 .bind(&[pull_id.into()])?
1218 .first::<LandRow>(None)
1219 .await
1220 }
1221
1222 async fn forget_landing(&self, pull_id: &str) -> Result<()> {
1223 self.db
1224 .prepare(
1225 "UPDATE pulls SET land_requested = NULL, land_requested_at = NULL WHERE id = ?",
1226 )
1227 .bind(&[pull_id.into()])?
1228 .run()
1229 .await?;
1230 Ok(())
1231 }
1232
1233 /// Lands a pull request whose head has just moved, if a merge of it was
1234 /// waiting for exactly that. What it was caught up to was already
1235 /// checked and reviewed apart from the merge, so the checks are not
1236 /// waited for again; a repository that wants them rerun turns on
1237 /// "require up to date", and then nothing is landed this way.
1238 pub(crate) async fn land_if_requested(&self, pull_id: &str) -> Result<()> {
1239 let Some(request) = self.land_request(pull_id).await? else {
1240 return Ok(());
1241 };
1242 self.forget_landing(pull_id).await?;
1243 let Some(pull) = self.pull_by_id(pull_id).await? else {
1244 return Ok(());
1245 };
1246 let repo: Outcome<Repo> = g1t_kit::call(
1247 &self.repos,
1248 "get_by_id",
1249 &GetByIdArgs {
1250 id: pull.repo_id.clone(),
1251 viewer: Some(request.actor.clone()),
1252 },
1253 )
1254 .await?;
1255 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
1256 return Ok(());
1257 };
1258 let merged = self
1259 .merge_pull(PullActionArgs {
1260 actor: request.actor,
1261 repo: RepoPath {
1262 namespace: repo.namespace,
1263 name: repo.name,
1264 },
1265 number: pull.number,
1266 summary: String::new(),
1267 keep_issue_open: request.keep_issue_open,
1268 ignore_checks: true,
1269 })
1270 .await?;
1271 if let Outcome::Fail(failure) = merged {
1272 self.stall(StallArgs {
1273 pull_id: pull.id,
1274 reason: format!(
1275 "It was brought up to date but could not be merged: {}",
1276 failure.message
1277 ),
1278 by: None,
1279 })
1280 .await?;
1281 }
1282 Ok(())
1283 }
1284
1285 /// What the runner needs to bring a pull request up to date for a merge
1286 /// that is waiting on it.
1287 pub(crate) async fn catch_up_job(&self, a: CatchUpJobArgs) -> Result<Option<LifecycleJob>> {
1288 if self.land_request(&a.pull_id).await?.is_none() {
1289 return Ok(None);
1290 }
1291 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
1292 return Ok(None);
1293 };
1294 // Its owner (whoever asked g1t for it, or its author) can read both
1295 // the repository and the pull request's source.
1296 let repo: Outcome<Repo> = g1t_kit::call(
1297 &self.repos,
1298 "get_by_id",
1299 &GetByIdArgs {
1300 id: pull.repo_id.clone(),
1301 viewer: self.owner_viewer(&pull).await?,
1302 },
1303 )
1304 .await?;
1305 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
1306 return Ok(None);
1307 };
1308 let path = RepoPath {
1309 namespace: repo.namespace,
1310 name: repo.name,
1311 };
1312 let issue = match pull.issue {
1313 Some(number) => self.issue(&pull.repo_id, number).await?,
1314 None => None,
1315 };
1316 let base = pull.base_branch(&repo.default_branch).to_owned();
1317 let feedback = self.conflict_note(&pull, &base).await?;
1318 Ok(Some(LifecycleJob {
1319 pull_id: pull.id,
1320 source: pull.fork.unwrap_or_else(|| path.clone()),
1321 repo: path,
1322 number: pull.number,
1323 author: pull.requested_by.unwrap_or(pull.author),
1324 branch: pull.branch,
1325 default_branch: base,
1326 title: pull.title,
1327 description: pull.body.unwrap_or_default(),
1328 issue,
1329 feedback,
1330 round: 0,
1331 }))
1332 }
1333
1334 /// For an agent catching up: the files g1t already knows conflict, so
1335 /// it reads them first. Empty when none are known.
1336 pub(crate) async fn conflict_note(&self, pull: &Pull, default_branch: &str) -> Result<String> {
1337 Ok(match self.conflicting_files(pull).await? {
1338 Some(files) if !files.is_empty() => format!(
1339 "g1t found ahead of time that merging {default_branch} into this pull request conflicts in these files: {}.",
1340 files.join(", ")
1341 ),
1342 _ => String::new(),
1343 })
1344 }
1345
1346 /// Stops seeing a pull request through until a person steps in. The
1347 /// first stop is published (`pull.stalled`), which tells its people
1348 /// that it needs them; a stop on one already stopped only says why.
1349 pub(crate) async fn stall(&self, a: StallArgs) -> Result<bool> {
1350 let before = self.pull_by_id(&a.pull_id).await?;
1351 let was_stalled = self.is_stalled(&a.pull_id).await?;
1352 self.db
1353 .prepare(
1354 "UPDATE pulls
1355 SET stalled = ?, working_on = NULL, working_until = NULL,
1356 land_requested = NULL, land_requested_at = NULL
1357 WHERE id = ? AND status = 'open'",
1358 )
1359 .bind(&[a.reason.trim().into(), a.pull_id.as_str().into()])?
1360 .run()
1361 .await?;
1362 self.db
1363 .prepare("UPDATE pulls SET stage = 'needs_you', stage_detail = ? WHERE id = ?")
1364 .bind(&[a.reason.trim().into(), a.pull_id.as_str().into()])?
1365 .run()
1366 .await?;
1367 if let Some(pull) = before.filter(|pull| !was_stalled && pull.status == PullStatus::Open) {
1368 self.publish_as(
1369 "pull.stalled",
1370 &pull.repo_id,
1371 a.by.clone(),
1372 g1t_contracts::events::PullEvent {
1373 detail: Some(a.reason.trim().to_owned()),
1374 ..Self::pull_event(&pull)
1375 },
1376 )
1377 .await?;
1378 }
1379 Ok(true)
1380 }
1381
1382 /// Whether g1t has stopped seeing the pull request through.
1383 pub(crate) async fn is_stalled(&self, pull_id: &str) -> Result<bool> {
1384 Ok(self
1385 .db
1386 .prepare("SELECT 1 AS value FROM pulls WHERE id = ? AND stalled IS NOT NULL")
1387 .bind(&[pull_id.into()])?
1388 .first::<u32>(Some("value"))
1389 .await?
1390 .is_some())
1391 }
1392
1393 /// Says that a pull request g1t had stopped on is going again, which
1394 /// closes what it was waiting on a person for.
1395 pub(crate) async fn announce_resumed(&self, pull_id: &str, actor: Option<String>) -> Result<()> {
1396 if let Some(pull) = self.pull_by_id(pull_id).await? {
1397 self.publish_as("pull.resumed", &pull.repo_id, actor, Self::pull_event(&pull)).await?;
1398 }
1399 Ok(())
1400 }
1401
1402 pub(crate) async fn managed_pulls(&self, a: ManagedPullsArgs) -> Result<Vec<String>> {
1403 let rows = self
1404 .db
1405 .prepare(
1406 "SELECT id AS value FROM pulls
1407 WHERE status = 'open' AND managed = 1 AND stalled IS NULL
1408 AND (?1 IS NULL OR repo_id = ?1)
1409 ORDER BY updated_at DESC LIMIT ?2",
1410 )
1411 .bind(&[
1412 a.repo_id.map_or(JsValue::NULL, JsValue::from),
1413 MANAGED_PAGE.into(),
1414 ])?
1415 .all()
1416 .await?
1417 .results::<ValueRow>()?;
1418 Ok(rows.into_iter().map(|row| row.value).collect())
1419 }
1420}
1421
1422#[cfg(test)]
1423mod tests {
1424 use super::*;
1425
1426 const MAX_REVISIONS: u32 = 2;
1427
1428 #[test]
1429 fn a_verdict_from_whoever_asked_for_g1t_s_change_is_not_someone_else_s() {
1430 use crate::rows::stored::{ASKER, G1T, pull};
1431 let made = pull(G1T, Some(ASKER));
1432 assert!(made_by_g1t(&made));
1433 // The person it was made for is held to what an author was: their
1434 // request for changes does not send g1t back as a reviewer's
1435 // would, and their approval does not lift a hold.
1436 assert!(!from_someone_else(&made, ASKER.0));
1437 assert!(!from_someone_else(&made, AGENT_ID));
1438 assert!(from_someone_else(&made, "usr_reviewer"));
1439 // Anyone's own pull request, the same.
1440 let own = pull(ASKER, None);
1441 assert!(!from_someone_else(&own, ASKER.0));
1442 assert!(from_someone_else(&own, "usr_reviewer"));
1443 }
1444
1445 #[test]
1446 fn approvals_the_repository_wants_are_waited_for() {
1447 let short = || Facts {
1448 review: reviewed(Some(Verdict::Approve)),
1449 approvals_missing: Some("This repository requires 1 approving review.".to_owned()),
1450 ..facts()
1451 };
1452 let (lifecycle, next) = decide(short());
1453 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1454 assert_eq!(
1455 lifecycle.detail,
1456 "This repository requires 1 approving review."
1457 );
1458 assert!(matches!(next, Next::Wait));
1459 // Not even a repository that merges by itself merges without them.
1460 let automatic = Facts {
1461 auto_merge: true,
1462 ..short()
1463 };
1464 assert_eq!(outcome(automatic), (Stage::NeedsYou, "wait"));
1465 }
1466
1467 #[test]
1468 fn a_repository_can_leave_review_to_people() {
1469 let unreviewed = Facts {
1470 agent_review: false,
1471 ..facts()
1472 };
1473 assert_eq!(outcome(unreviewed), (Stage::Ready, "wait"));
1474 let failing = Facts {
1475 agent_review: false,
1476 check_status: Some(CheckStatus::Failed),
1477 ..facts()
1478 };
1479 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1480 }
1481
1482 #[test]
1483 fn a_repository_sets_how_often_the_author_is_sent_back() {
1484 let never = Facts {
1485 max_revisions: 0,
1486 check_status: Some(CheckStatus::Failed),
1487 ..facts()
1488 };
1489 assert_eq!(outcome(never), (Stage::NeedsYou, "wait"));
1490 }
1491
1492 /// Checks on a head commit: `(context, state)` statuses, against the
1493 /// required check names.
1494 fn checks(statuses: &[(&str, &str)], required: &[&str]) -> WorkflowFacts {
1495 let statuses: Vec<CommitStatus> = statuses
1496 .iter()
1497 .map(|(context, state)| CommitStatus {
1498 context: (*context).to_owned(),
1499 state: (*state).to_owned(),
1500 description: None,
1501 target_url: None,
1502 updated_at: String::new(),
1503 })
1504 .collect();
1505 let required: Vec<String> = required.iter().map(|name| (*name).to_owned()).collect();
1506 WorkflowFacts::of(&statuses, &required)
1507 }
1508
1509 /// A pull request that is ready for review, whose required check
1510 /// passed, and nothing else yet.
1511 fn facts() -> Facts {
1512 Facts {
1513 draft: false,
1514 stalled: None,
1515 working_on: None,
1516 check_status: None,
1517 review_pending: false,
1518 revisions: 0,
1519 review: None,
1520 behind: false,
1521 conflicting: false,
1522 auto_merge: false,
1523 require_up_to_date: false,
1524 agent_review: true,
1525 max_revisions: MAX_REVISIONS,
1526 approvals_missing: None,
1527 person_request: None,
1528 queued: None,
1529 workflows: checks(&[("CI / pull_request", "success")], &["CI"]),
1530 low_confidence: None,
1531 }
1532 }
1533
1534 #[test]
1535 fn failed_checks_send_the_agent_back_and_running_ones_wait() {
1536 let failed = Facts {
1537 workflows: checks(&[("CI / pull_request", "failure")], &["CI"]),
1538 ..facts()
1539 };
1540 let (lifecycle, next) = decide(failed);
1541 assert!(matches!(next, Next::Revise(Feedback::FailedWorkflows)));
1542 assert!(lifecycle.detail.contains("CI / pull_request failed"));
1543 let running = Facts {
1544 workflows: checks(&[("CI / pull_request", "pending")], &["CI"]),
1545 ..facts()
1546 };
1547 let (lifecycle, next) = decide(running);
1548 assert!(matches!(next, Next::Wait));
1549 assert!(lifecycle.detail.contains("Waiting for CI / pull_request"));
1550 }
1551
1552 #[test]
1553 fn a_check_the_branch_does_not_require_is_fixed_but_does_not_hold_it_for_ever() {
1554 // Lint is not required: the agent is still sent back to fix it...
1555 let lint = checks(&[("CI / pull_request", "success"), ("Lint / pull_request", "failure")], &["CI"]);
1556 let failing = Facts { workflows: lint.clone(), ..facts() };
1557 assert_eq!(outcome(failing), (Stage::Revising, "revise for workflows"));
1558 // ...but once it is out of revisions, Lint no longer holds it.
1559 let exhausted = Facts {
1560 workflows: lint,
1561 revisions: MAX_REVISIONS,
1562 review: reviewed(Some(Verdict::Approve)),
1563 ..facts()
1564 };
1565 assert_eq!(outcome(exhausted), (Stage::Ready, "wait"));
1566 // A required check that still fails asks a person.
1567 let required = Facts {
1568 workflows: checks(&[("CI / pull_request", "failure")], &["CI"]),
1569 revisions: MAX_REVISIONS,
1570 ..facts()
1571 };
1572 let (lifecycle, next) = decide(required);
1573 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1574 assert_eq!(lifecycle.detail, "The required check CI still fails after the agent revised twice.");
1575 assert!(matches!(next, Next::Wait));
1576 }
1577
1578 #[test]
1579 fn a_required_check_that_has_not_reported_is_waited_for() {
1580 let missing = Facts {
1581 workflows: checks(&[("CI / pull_request", "success")], &["CI", "Deploy"]),
1582 review: reviewed(Some(Verdict::Approve)),
1583 auto_merge: true,
1584 ..facts()
1585 };
1586 let (lifecycle, next) = decide(missing);
1587 assert_eq!(lifecycle.stage, Stage::Checking);
1588 assert_eq!(lifecycle.detail, "Waiting for the required check Deploy to report on its latest commit.");
1589 assert!(matches!(next, Next::Wait), "auto-merge must not land it");
1590 }
1591
1592 #[test]
1593 fn a_queued_pull_request_waits_in_the_queue() {
1594 let queued = Facts {
1595 review: reviewed(Some(Verdict::Approve)),
1596 queued: Some((QueueState::Testing, vec![12, 14])),
1597 auto_merge: true,
1598 ..facts()
1599 };
1600 let (lifecycle, next) = decide(queued);
1601 assert_eq!(lifecycle.stage, Stage::Queued);
1602 assert!(lifecycle.detail.contains("#12, #14"));
1603 assert!(matches!(next, Next::Wait));
1604 }
1605
1606 fn asked_by_a_person() -> Option<PersonRequest> {
1607 Some(PersonRequest {
1608 author_id: "usr_reviewer".to_owned(),
1609 author_name: "g1t-reviewer".to_owned(),
1610 created_at: "2026-10-02T11:00:00.000Z".to_owned(),
1611 })
1612 }
1613
1614 #[test]
1615 fn a_person_asking_for_changes_sends_the_agent_back() {
1616 let asked = Facts {
1617 review: reviewed(Some(Verdict::Approve)),
1618 approvals_missing: Some("A reviewer has asked for changes.".to_owned()),
1619 person_request: asked_by_a_person(),
1620 ..facts()
1621 };
1622 let (lifecycle, _) = decide(Facts {
1623 person_request: asked_by_a_person(),
1624 ..facts()
1625 });
1626 assert!(lifecycle.detail.starts_with("g1t-reviewer asked for changes"));
1627 assert_eq!(outcome(asked), (Stage::Revising, "revise for a person"));
1628 }
1629
1630 #[test]
1631 fn a_person_is_asked_once_the_revisions_run_out() {
1632 let exhausted = Facts {
1633 person_request: asked_by_a_person(),
1634 revisions: MAX_REVISIONS,
1635 ..facts()
1636 };
1637 assert_eq!(outcome(exhausted), (Stage::NeedsYou, "wait"));
1638 }
1639
1640 fn reviewed(verdict: Option<Verdict>) -> Option<FinishedReview> {
1641 Some(FinishedReview {
1642 finished_at: "2026-10-02T10:00:00.000Z".to_owned(),
1643 verdict,
1644 })
1645 }
1646
1647 /// The stage, and a word for the step to take.
1648 fn outcome(facts: Facts) -> (Stage, &'static str) {
1649 let (lifecycle, next) = decide(facts);
1650 let step = match next {
1651 Next::Wait => "wait",
1652 Next::Review => "review",
1653 Next::Revise(Feedback::FailedChecks) => "revise for checks",
1654 Next::Revise(Feedback::FailedWorkflows) => "revise for workflows",
1655 Next::Revise(Feedback::Review(_)) => "revise for review",
1656 Next::Revise(Feedback::Person(_)) => "revise for a person",
1657 Next::CatchUp => "catch up",
1658 Next::Merge => "merge",
1659 };
1660 (lifecycle.stage, step)
1661 }
1662
1663 #[test]
1664 fn nothing_is_started_while_the_agent_is_still_working() {
1665 let draft = Facts {
1666 draft: true,
1667 check_status: None,
1668 ..facts()
1669 };
1670 assert_eq!(outcome(draft), (Stage::Working, "wait"));
1671 }
1672
1673 #[test]
1674 fn checks_come_before_review() {
1675 let unchecked = Facts {
1676 workflows: checks(&[], &["CI"]),
1677 ..facts()
1678 };
1679 assert_eq!(outcome(unchecked), (Stage::Checking, "wait"));
1680 let running = Facts {
1681 workflows: checks(&[("CI / pull_request", "pending")], &["CI"]),
1682 ..facts()
1683 };
1684 assert_eq!(outcome(running), (Stage::Checking, "wait"));
1685 assert_eq!(outcome(facts()), (Stage::Reviewing, "review"));
1686 }
1687
1688 #[test]
1689 fn a_branch_that_requires_no_checks_goes_straight_to_review() {
1690 let unchecked = Facts {
1691 workflows: WorkflowFacts::default(),
1692 ..facts()
1693 };
1694 assert_eq!(outcome(unchecked), (Stage::Reviewing, "review"));
1695 }
1696
1697 #[test]
1698 fn failing_in_the_merge_queue_sends_the_author_back() {
1699 let failed = Facts {
1700 check_status: Some(CheckStatus::Failed),
1701 ..facts()
1702 };
1703 let (lifecycle, next) = decide(failed);
1704 assert_eq!(lifecycle.detail, "It failed in the merge queue. The agent is being sent back to fix it.");
1705 assert!(matches!(next, Next::Revise(Feedback::FailedChecks)));
1706 }
1707
1708 #[test]
1709 fn checks_that_could_not_run_are_a_persons_problem() {
1710 let errored = Facts {
1711 check_status: Some(CheckStatus::Errored),
1712 ..facts()
1713 };
1714 assert_eq!(outcome(errored), (Stage::NeedsYou, "wait"));
1715 }
1716
1717 #[test]
1718 fn a_review_asking_for_changes_sends_the_author_back() {
1719 let changes = Facts {
1720 review: reviewed(Some(Verdict::RequestChanges)),
1721 ..facts()
1722 };
1723 assert_eq!(outcome(changes), (Stage::Revising, "revise for review"));
1724 }
1725
1726 #[test]
1727 fn the_author_is_sent_back_only_so_many_times() {
1728 let failing = Facts {
1729 check_status: Some(CheckStatus::Failed),
1730 revisions: MAX_REVISIONS,
1731 ..facts()
1732 };
1733 assert_eq!(outcome(failing), (Stage::NeedsYou, "wait"));
1734 let unconvinced = Facts {
1735 review: reviewed(Some(Verdict::RequestChanges)),
1736 revisions: MAX_REVISIONS,
1737 ..facts()
1738 };
1739 assert_eq!(outcome(unconvinced), (Stage::NeedsYou, "wait"));
1740 // One short of the limit still gets another go.
1741 let once = Facts {
1742 check_status: Some(CheckStatus::Failed),
1743 revisions: MAX_REVISIONS - 1,
1744 ..facts()
1745 };
1746 assert_eq!(outcome(once), (Stage::Revising, "revise for checks"));
1747 }
1748
1749 #[test]
1750 fn a_review_that_could_not_be_written_is_not_retried() {
1751 let broken = Facts {
1752 review: reviewed(None),
1753 ..facts()
1754 };
1755 assert_eq!(outcome(broken), (Stage::NeedsYou, "wait"));
1756 }
1757
1758 #[test]
1759 fn being_behind_only_holds_a_change_up_where_the_repository_says_so() {
1760 let behind = || Facts {
1761 review: reviewed(Some(Verdict::Approve)),
1762 behind: true,
1763 ..facts()
1764 };
1765 // By default it is ready as it is; merging brings it up to date.
1766 assert_eq!(outcome(behind()), (Stage::Ready, "wait"));
1767 let strict = Facts {
1768 require_up_to_date: true,
1769 ..behind()
1770 };
1771 assert_eq!(outcome(strict), (Stage::CatchingUp, "catch up"));
1772 let current = Facts {
1773 review: reviewed(Some(Verdict::Approve)),
1774 ..facts()
1775 };
1776 assert_eq!(outcome(current), (Stage::Ready, "wait"));
1777 }
1778
1779 #[test]
1780 fn a_ready_change_lands_by_itself_only_where_the_repository_says_so() {
1781 let ready = || Facts {
1782 review: reviewed(Some(Verdict::Approve)),
1783 ..facts()
1784 };
1785 assert_eq!(outcome(ready()), (Stage::Ready, "wait"));
1786 let automatic = Facts {
1787 auto_merge: true,
1788 ..ready()
1789 };
1790 assert_eq!(outcome(automatic), (Stage::Ready, "merge"));
1791 // One that is behind is merged too: merging brings it up to date.
1792 let behind = Facts {
1793 auto_merge: true,
1794 behind: true,
1795 ..ready()
1796 };
1797 assert_eq!(outcome(behind), (Stage::Ready, "merge"));
1798 // Unless the repository wants it caught up and checked again first.
1799 let strict = Facts {
1800 auto_merge: true,
1801 behind: true,
1802 require_up_to_date: true,
1803 ..ready()
1804 };
1805 assert_eq!(outcome(strict), (Stage::CatchingUp, "catch up"));
1806 // Nothing short of approved is merged, whatever the setting.
1807 let failing = Facts {
1808 auto_merge: true,
1809 check_status: Some(CheckStatus::Failed),
1810 ..ready()
1811 };
1812 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1813 let unreviewed = Facts {
1814 auto_merge: true,
1815 ..facts()
1816 };
1817 assert_eq!(outcome(unreviewed), (Stage::Reviewing, "review"));
1818 }
1819
1820 #[test]
1821 fn a_conflict_found_ahead_of_time_is_resolved_before_merging() {
1822 let conflicting = || Facts {
1823 review: reviewed(Some(Verdict::Approve)),
1824 behind: true,
1825 conflicting: true,
1826 ..facts()
1827 };
1828 // Even where the repository would merge one that is merely behind.
1829 assert_eq!(outcome(conflicting()), (Stage::CatchingUp, "catch up"));
1830 let automatic = Facts {
1831 auto_merge: true,
1832 ..conflicting()
1833 };
1834 assert_eq!(outcome(automatic), (Stage::CatchingUp, "catch up"));
1835 // Failed checks come first: a revision merges the branch in too.
1836 let failing = Facts {
1837 check_status: Some(CheckStatus::Failed),
1838 ..conflicting()
1839 };
1840 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1841 }
1842
1843 #[test]
1844 fn catching_up_reruns_the_checks_but_not_the_review() {
1845 // The merge moved the head, so the checks are waited for again.
1846 let merged_in = Facts {
1847 review: reviewed(Some(Verdict::Approve)),
1848 workflows: checks(&[], &["CI"]),
1849 ..facts()
1850 };
1851 assert_eq!(outcome(merged_in), (Stage::Checking, "wait"));
1852 }
1853
1854 #[test]
1855 fn a_step_under_way_is_not_started_again() {
1856 for (step, stage) in [
1857 ("review", Stage::Reviewing),
1858 ("revision", Stage::Revising),
1859 ("catch_up", Stage::CatchingUp),
1860 ] {
1861 let busy = Facts {
1862 working_on: Some(step.to_owned()),
1863 // Whatever else is true, the step in hand comes first.
1864 check_status: Some(CheckStatus::Failed),
1865 ..facts()
1866 };
1867 assert_eq!(outcome(busy), (stage, "wait"));
1868 }
1869 let asked = Facts {
1870 review_pending: true,
1871 ..facts()
1872 };
1873 assert_eq!(outcome(asked), (Stage::Reviewing, "wait"));
1874 }
1875
1876 #[test]
1877 fn a_low_confidence_change_waits_for_a_person_instead_of_merging() {
1878 let held = || Facts {
1879 review: reviewed(Some(Verdict::Approve)),
1880 auto_merge: true,
1881 low_confidence: Some("tests not added, 3 revisions".to_owned()),
1882 ..facts()
1883 };
1884 let (lifecycle, next) = decide(held());
1885 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1886 assert!(matches!(next, Next::Wait), "auto-merge must not land it");
1887 assert_eq!(
1888 lifecycle.detail,
1889 "The agent's confidence in this change is low (tests not added, 3 revisions). This repository asks a person before merging it: approve it to let it land, or ask for changes."
1890 );
1891 // Without auto-merge it needs someone too, and says why.
1892 assert_eq!(outcome(Facts { auto_merge: false, ..held() }), (Stage::NeedsYou, "wait"));
1893 // Not held (the setting is off, or a person approved): it lands.
1894 assert_eq!(outcome(Facts { low_confidence: None, ..held() }), (Stage::Ready, "merge"));
1895 // It holds only a change that is otherwise ready: what comes first,
1896 // such as failed checks, is still dealt with first.
1897 let failing = Facts {
1898 check_status: Some(CheckStatus::Failed),
1899 ..held()
1900 };
1901 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1902 let unapproved = Facts {
1903 approvals_missing: Some("This repository requires 1 approving review.".to_owned()),
1904 ..held()
1905 };
1906 assert_eq!(decide(unapproved).0.detail, "This repository requires 1 approving review.");
1907 }
1908
1909 #[test]
1910 fn once_stopped_it_stays_stopped() {
1911 let (lifecycle, next) = decide(Facts {
1912 stalled: Some("The agent could not catch up.".to_owned()),
1913 review: reviewed(Some(Verdict::Approve)),
1914 behind: true,
1915 ..facts()
1916 });
1917 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1918 assert_eq!(lifecycle.detail, "The agent could not catch up.");
1919 assert!(matches!(next, Next::Wait));
1920 }
1921}