Skip to content

g1t/services/work/src/lifecycle.rs

1,942 lines75,552 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 // The rules of the branch it merges into, as g1t sees them when it
605 // merges by itself: what people must still do, whether an agent's
606 // change may land here unattended, and a cost cap that holds the
607 // agent until a person approves.
608 let gate = match self.repo_for(pull).await? {
609 Some(repo) => Some(self.merge_gate(&repo, pull, Some(&User::system(&repo.namespace)), false, true).await?),
610 None => None,
611 };
612 let level = confidence.as_ref().map(|confidence| confidence.level).or(pull.confidence.as_ref().map(|c| c.level));
613 let rules_allow_auto_merge = gate
614 .as_ref()
615 .is_none_or(|gate| g1t_rules::merge::auto_merge_refusal(&gate.requirements, level).is_none());
616 let held = gate.as_ref().and_then(|gate| {
617 g1t_rules::outcome::blocking(&gate.judged)
618 .into_iter()
619 .find(|violation| violation.rule == "cost_cap")
620 .map(|violation| format!("{} {}", violation.message, violation.remedy))
621 });
622 let (lifecycle, next) = decide(Facts {
623 draft: pull.status == PullStatus::Draft,
624 stalled: progress.stalled.or(held),
625 working_on,
626 check_status: pull.check_status,
627 review_pending: self.review_pending(&pull.id).await?,
628 revisions: progress.revisions,
629 review,
630 behind,
631 conflicting: self.conflicting_files(pull).await?.is_some(),
632 auto_merge: settings.auto_merge && rules_allow_auto_merge,
633 require_up_to_date: settings.require_up_to_date,
634 agent_review: settings.agent_review,
635 max_revisions: settings.max_revisions,
636 approvals_missing: gate.as_ref().and_then(|gate| gate.people_gap()),
637 person_request: self
638 .person_request(pull, progress.revised_at.as_deref())
639 .await?,
640 queued: self.queued_entry(&pull.id).await?,
641 workflows,
642 low_confidence,
643 });
644 Ok(Some((lifecycle, next, confidence)))
645 }
646
647 /// The latest request for changes by a person other than its owner
648 /// (whoever asked g1t for it, or its author), if it is that person's
649 /// latest verdict and came after the last revision.
650 async fn person_request(
651 &self,
652 pull: &Pull,
653 revised_at: Option<&str>,
654 ) -> Result<Option<PersonRequest>> {
655 #[derive(Deserialize)]
656 struct Verdicts {
657 author_id: String,
658 author_name: String,
659 verdict: Verdict,
660 created_at: String,
661 }
662 let rows = match self.prefetched_pull(&pull.id) {
663 Some(found) => found
664 .rows::<Verdicts>(Slot::Verdicts)?
665 .into_iter()
666 .filter(|row| from_someone_else(pull, &row.author_id))
667 .collect(),
668 None => self
669 .db
670 .prepare(
671 "SELECT author_id, author_name, verdict, created_at FROM comments
672 WHERE repo_id = ? AND number = ? AND verdict IS NOT NULL
673 AND author_id != ? AND author_id != ?
674 ORDER BY id",
675 )
676 .bind(&[
677 pull.repo_id.as_str().into(),
678 pull.number.into(),
679 pull.owner().id.as_str().into(),
680 AGENT_ID.into(),
681 ])?
682 .all()
683 .await?
684 .results::<Verdicts>()?,
685 };
686 // Each person's latest verdict is the one that stands.
687 let mut latest: HashMap<String, Verdicts> = HashMap::new();
688 for row in rows {
689 latest.insert(row.author_id.clone(), row);
690 }
691 Ok(latest
692 .into_values()
693 .filter(|row| row.verdict == Verdict::RequestChanges)
694 .filter(|row| revised_at.is_none_or(|revised| row.created_at.as_str() > revised))
695 .max_by(|a, b| a.created_at.cmp(&b.created_at))
696 .map(|row| PersonRequest {
697 author_id: row.author_id,
698 author_name: row.author_name,
699 created_at: row.created_at,
700 }))
701 }
702
703 /// Marks a pull request a g1t agent has just opened as one g1t sees
704 /// through.
705 pub(crate) async fn manage(&self, pull: &Pull) -> Result<()> {
706 if !made_by_g1t(pull) {
707 return Ok(());
708 }
709 self.db
710 .prepare("UPDATE pulls SET managed = 1 WHERE id = ?")
711 .bind(&[pull.id.as_str().into()])?
712 .run()
713 .await?;
714 Ok(())
715 }
716
717 /// Takes a step for a pull request, if nobody else has. One statement,
718 /// so that two callers cannot both take it.
719 pub(crate) async fn claim(&self, pull_id: &str, step: &str, minutes: u64, revising: bool) -> Result<bool> {
720 let now = now_ms();
721 let revision = if revising {
722 ", revisions = revisions + 1, revised_at = ?1"
723 } else {
724 ""
725 };
726 Ok(self
727 .db
728 .prepare(format!(
729 "UPDATE pulls SET working_on = ?2, working_until = ?3{revision}
730 WHERE id = ?4 AND status = 'open' AND stalled IS NULL
731 AND (working_until IS NULL OR working_until < ?1)
732 RETURNING id AS value"
733 ))
734 .bind(&[
735 rfc3339(now).into(),
736 step.into(),
737 rfc3339(now + minutes * 60 * 1000).into(),
738 pull_id.into(),
739 ])?
740 .first::<ValueRow>(None)
741 .await?
742 .is_some())
743 }
744
745 /// What the author is told when sent back: the checks that failed and
746 /// what their failing jobs printed, why the merge queue took it out, or
747 /// the review and its comments on lines.
748 async fn feedback(&self, pull: &Pull, feedback: &Feedback) -> Result<String> {
749 match feedback {
750 Feedback::FailedChecks => {
751 let run = self.latest_checks(&pull.id).await?;
752 // Why it failed, when that is more than a list of commands:
753 // the merge queue saying what broke in the combined state.
754 let why = run.as_ref().and_then(|run| run.error.clone()).map(|error| format!("{error}\n\n")).unwrap_or_default();
755 let failed: Vec<String> = run
756 .map(|run| run.results)
757 .unwrap_or_default()
758 .into_iter()
759 .filter(|result| !result.passed)
760 .map(|result| {
761 let length = result.output.chars().count();
762 let output: String = result
763 .output
764 .chars()
765 .skip(length.saturating_sub(MAX_CHECK_OUTPUT_CHARS))
766 .collect();
767 let exit = result
768 .exit_code
769 .map_or("it was stopped for taking too long".to_owned(), |code| {
770 format!("exit code {code}")
771 });
772 format!("`{}` failed ({exit}):\n\n{}", result.command, output.trim())
773 })
774 .collect();
775 if failed.is_empty() {
776 return Ok(format!(
777 "{why}Find the cause, fix it in your change, and push. Use get_workflow_run and get_job_logs for any workflow named above."
778 ));
779 }
780 Ok(format!(
781 "{why}These commands failed against your change.\n\n{}",
782 failed.join("\n\n")
783 ))
784 }
785 Feedback::FailedWorkflows => {
786 let (statuses, settings) = futures_util::future::try_join(
787 self.statuses(&pull.repo_id, pull.head_commit.as_deref()),
788 self.settings_for(pull),
789 )
790 .await?;
791 let required = |context: &str| {
792 let name = g1t_contracts::work::check_name(context).0;
793 settings.required_checks.iter().any(|wanted| wanted.eq_ignore_ascii_case(name))
794 };
795 let failing: Vec<&CommitStatus> =
796 statuses.iter().filter(|s| s.state == "failure" || s.state == "error").collect();
797 let failed: Vec<String> = failing
798 .iter()
799 .map(|s| {
800 let run = s.target_url.as_deref().and_then(|url| url.rsplit('/').next()).unwrap_or_default();
801 format!(
802 "- {}{} ({}): run `{run}`, {}",
803 s.context,
804 if required(&s.context) { ", required to merge" } else { "" },
805 s.description.as_deref().unwrap_or("failed"),
806 s.target_url.as_deref().unwrap_or_default()
807 )
808 })
809 .collect();
810 let logs = self.failing_logs(pull, &failing).await.unwrap_or_default();
811 // Named outright: the agent cannot guess it from its fork.
812 let repo = g1t_kit::call::<_, Option<RepoPath>>(
813 &self.repos,
814 "path_by_id",
815 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
816 )
817 .await?
818 .map(|path| format!("{}/{}", path.namespace, path.name))
819 .unwrap_or_default();
820 let logs = if logs.is_empty() {
821 String::new()
822 } else {
823 format!("\n\nThe end of what the failing jobs printed:\n\n{}", logs.join("\n\n"))
824 };
825 Ok(format!(
826 "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\
827 For more, use the `get_workflow_run` tool (repo `{repo}` and the run's id), \
828 then `get_job_logs` for the job that failed. Fix the cause in the code, not the workflow, \
829 unless the workflow itself is wrong. Push, and the checks run again.",
830 failed.join("\n")
831 ))
832 }
833 Feedback::Review(finished_at) => {
834 // Everything a review says is recorded at the moment it finished.
835 let notes = self
836 .db
837 .prepare(
838 "SELECT body, path, line FROM comments
839 WHERE repo_id = ? AND number = ? AND author_id = ? AND created_at = ?
840 ORDER BY id",
841 )
842 .bind(&[
843 pull.repo_id.as_str().into(),
844 pull.number.into(),
845 AGENT_ID.into(),
846 finished_at.as_str().into(),
847 ])?
848 .all()
849 .await?
850 .results::<ReviewNote>()?;
851 let mut on_lines = Vec::new();
852 let mut summary = String::new();
853 for note in notes {
854 match (note.path, note.line) {
855 (Some(path), Some(line)) => {
856 on_lines.push(format!("- `{path}` line {line}: {}", note.body));
857 }
858 (Some(path), None) => on_lines.push(format!("- `{path}`: {}", note.body)),
859 (None, _) => summary = note.body,
860 }
861 }
862 let mut text = format!(
863 "Another agent reviewed your change and asked for changes.\n\n{summary}"
864 );
865 if !on_lines.is_empty() {
866 text.push_str("\n\nIts comments on lines:\n");
867 text.push_str(&on_lines.join("\n"));
868 }
869 Ok(text)
870 }
871 Feedback::Person(request) => {
872 // What they wrote since the agent last revised, which their
873 // request for changes closes.
874 let revised: Option<String> = self
875 .db
876 .prepare("SELECT revised_at AS value FROM pulls WHERE id = ?")
877 .bind(&[pull.id.as_str().into()])?
878 .first::<Option<String>>(Some("value"))
879 .await?
880 .flatten();
881 let notes = self
882 .db
883 .prepare(
884 "SELECT body, path, line FROM comments
885 WHERE repo_id = ? AND number = ? AND author_id = ?
886 AND created_at > ? AND created_at <= ?
887 ORDER BY id",
888 )
889 .bind(&[
890 pull.repo_id.as_str().into(),
891 pull.number.into(),
892 request.author_id.as_str().into(),
893 revised.unwrap_or_default().into(),
894 request.created_at.as_str().into(),
895 ])?
896 .all()
897 .await?
898 .results::<ReviewNote>()?;
899 let mut on_lines = Vec::new();
900 let mut said = Vec::new();
901 for note in notes {
902 match (note.path, note.line) {
903 (Some(path), Some(line)) => {
904 on_lines.push(format!("- `{path}` line {line}: {}", note.body));
905 }
906 (Some(path), None) => on_lines.push(format!("- `{path}`: {}", note.body)),
907 (None, _) => said.push(note.body),
908 }
909 }
910 let mut text = format!(
911 "{} reviewed your change and asked for changes.\n\n{}",
912 request.author_name,
913 said.join("\n\n")
914 );
915 if !on_lines.is_empty() {
916 text.push_str("\n\nTheir comments on lines:\n");
917 text.push_str(&on_lines.join("\n"));
918 }
919 Ok(text)
920 }
921 }
922 }
923
924 /// The end of what the failed jobs of failing workflow runs printed, a
925 /// few jobs at most, for an agent sent back to fix them. Empty where the
926 /// runs or their logs cannot be read.
927 async fn failing_logs(&self, pull: &Pull, failing: &[&CommitStatus]) -> Result<Vec<String>> {
928 use g1t_contracts::actions::{JobLog, LogsArgs, RunArgs, RunDetail};
929 let Some(repo) = g1t_kit::call::<_, Option<RepoPath>>(
930 &self.repos,
931 "path_by_id",
932 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
933 )
934 .await?
935 else {
936 return Ok(Vec::new());
937 };
938 let viewer = self.owner_viewer(pull).await?;
939 let mut out = Vec::new();
940 for status in failing {
941 let Some(run_id) = status
942 .target_url
943 .as_deref()
944 .and_then(|url| url.split("/actions/runs/").nth(1))
945 .map(|rest| rest.split(['/', '?', '#']).next().unwrap_or_default().to_owned())
946 .filter(|id| !id.is_empty())
947 else {
948 continue;
949 };
950 let detail: Outcome<RunDetail> = g1t_kit::call(
951 &self.actions,
952 "run",
953 &RunArgs { repo: repo.clone(), viewer: viewer.clone(), id: run_id },
954 )
955 .await?;
956 let Outcome::Ok(detail) = detail else { continue };
957 let failed_jobs = detail.jobs.into_iter().filter(|job| {
958 job.status == "completed" && matches!(job.conclusion.as_deref(), Some("failure" | "timed_out"))
959 });
960 for job in failed_jobs {
961 if out.len() >= MAX_FAILED_JOBS {
962 return Ok(out);
963 }
964 let mut text = String::new();
965 let mut after = 0;
966 for _ in 0..MAX_LOG_PAGES {
967 let page: Outcome<JobLog> = g1t_kit::call(
968 &self.actions,
969 "logs",
970 &LogsArgs { repo: repo.clone(), viewer: viewer.clone(), job: job.id.clone(), after },
971 )
972 .await?;
973 let Outcome::Ok(page) = page else { break };
974 let Some(last) = page.chunks.last().map(|chunk| chunk.seq) else { break };
975 for chunk in &page.chunks {
976 text.push_str(&chunk.text);
977 if !chunk.text.ends_with('\n') {
978 text.push('\n');
979 }
980 }
981 // Only the end is kept, so the start can go as it is read.
982 let length = text.chars().count();
983 if length > MAX_JOB_LOG_CHARS * 2 {
984 text = text.chars().skip(length - MAX_JOB_LOG_CHARS).collect();
985 }
986 after = last;
987 if page.chunks.len() < 500 {
988 break;
989 }
990 }
991 let length = text.chars().count();
992 let tail: String = text.chars().skip(length.saturating_sub(MAX_JOB_LOG_CHARS)).collect();
993 if !tail.trim().is_empty() {
994 out.push(format!("{} / {}:\n```\n{}\n```", status.context, job.name, tail.trim_end()));
995 }
996 }
997 }
998 Ok(out)
999 }
1000
1001 pub(crate) async fn advance(&self, a: AdvanceArgs) -> Result<Advance> {
1002 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
1003 return Ok(Advance::None);
1004 };
1005 if pull.status != PullStatus::Open {
1006 return Ok(Advance::None);
1007 }
1008 // As a member: a private repository would look missing otherwise,
1009 // and the pull request would never move.
1010 let viewer: Viewer = self.owner_viewer(&pull).await?;
1011 let repo: Outcome<Repo> = g1t_kit::call(
1012 &self.repos,
1013 "get_by_id",
1014 &GetByIdArgs {
1015 id: pull.repo_id.clone(),
1016 viewer,
1017 },
1018 )
1019 .await?;
1020 let (Outcome::Ok(repo), Some(source)) = (crate::retired::unless_archived(repo), pull.fork.clone()) else {
1021 return Ok(Advance::None);
1022 };
1023 // The branch it merges into: what it catches up with.
1024 let base = pull.base_branch(&repo.default_branch).to_owned();
1025 let issue = match pull.issue {
1026 Some(number) => self.issue(&pull.repo_id, number).await?,
1027 None => None,
1028 };
1029 let behind = self.is_behind(&repo.id, &pull).await?;
1030 let Some((lifecycle, next)) = self.assess(&pull, &issue, behind).await? else {
1031 return Ok(Advance::None);
1032 };
1033
1034 if matches!(next, Next::Merge) {
1035 self.merge_by_policy(&repo, &pull).await?;
1036 return Ok(Advance::None);
1037 }
1038 let (step, minutes) = match &next {
1039 Next::Wait | Next::Merge => return Ok(Advance::None),
1040 Next::Review => ("review", REVIEW_MINUTES),
1041 Next::Revise(_) => ("revision", REVISION_MINUTES),
1042 Next::CatchUp => ("catch_up", CATCH_UP_MINUTES),
1043 };
1044 let feedback = match &next {
1045 Next::Revise(feedback) => self.feedback(&pull, feedback).await?,
1046 Next::CatchUp => self.conflict_note(&pull, &base).await?,
1047 _ => String::new(),
1048 };
1049 if !self
1050 .claim(&pull.id, step, minutes, matches!(next, Next::Revise(_)))
1051 .await?
1052 {
1053 return Ok(Advance::None);
1054 }
1055 // Said in the conversation, so nobody has to wonder why a review or
1056 // a new commit appeared.
1057 let told = match &next {
1058 Next::Review => {
1059 self.db
1060 .prepare(
1061 "UPDATE pulls SET reviewers = json_insert(reviewers, '$[#]', ?1)
1062 WHERE id = ?2 AND NOT EXISTS (
1063 SELECT 1 FROM json_each(pulls.reviewers) WHERE json_each.value = ?1)",
1064 )
1065 .bind(&[AGENT_NAME.into(), pull.id.as_str().into()])?
1066 .run()
1067 .await?;
1068 "requested a review from g1t".to_owned()
1069 }
1070 Next::Revise(Feedback::FailedChecks) => {
1071 "sent g1t back to fix what failed in the merge queue".to_owned()
1072 }
1073 Next::Revise(Feedback::FailedWorkflows) => {
1074 "sent g1t back to fix the failed checks".to_owned()
1075 }
1076 Next::Revise(Feedback::Review(_)) => {
1077 "sent g1t back to address the review".to_owned()
1078 }
1079 _ => format!("asked g1t to bring this up to date with {base}"),
1080 };
1081 self.note(
1082 &pull.repo_id,
1083 pull.number,
1084 (POLICY_ACTOR_ID, POLICY_ACTOR_NAME),
1085 &told,
1086 )
1087 .await?;
1088 let job = LifecycleJob {
1089 pull_id: pull.id,
1090 repo: RepoPath {
1091 namespace: repo.namespace,
1092 name: repo.name,
1093 },
1094 number: pull.number,
1095 author: pull.requested_by.unwrap_or(pull.author),
1096 source,
1097 branch: None,
1098 default_branch: base,
1099 title: pull.title,
1100 description: pull.body.unwrap_or_default(),
1101 issue,
1102 feedback,
1103 round: lifecycle.revisions + 1,
1104 };
1105 Ok(match next {
1106 Next::Review => Advance::Review { job },
1107 Next::Revise(_) => Advance::Revise { job },
1108 Next::CatchUp => Advance::CatchUp { job },
1109 Next::Wait | Next::Merge => Advance::None,
1110 })
1111 }
1112
1113 /// Lands a pull request that is ready, on the authority of the
1114 /// repository's settings instead of a person's click.
1115 async fn merge_by_policy(&self, repo: &Repo, pull: &Pull) -> Result<()> {
1116 if !self.claim(&pull.id, "merge", MERGE_MINUTES, false).await? {
1117 return Ok(());
1118 }
1119 // g1t acts for the workspace whose members turned this on, with the
1120 // default base permission (Write), which merging needs.
1121 let actor = User {
1122 id: POLICY_ACTOR_ID.to_owned(),
1123 username: POLICY_ACTOR_NAME.to_owned(),
1124 verified: true,
1125 workspaces: vec![Membership::member(repo.namespace.to_lowercase())],
1126 ..User::default()
1127 };
1128 let merged = self
1129 .merge_pull(PullActionArgs {
1130 actor,
1131 repo: RepoPath {
1132 namespace: repo.namespace.clone(),
1133 name: repo.name.clone(),
1134 },
1135 number: pull.number,
1136 summary: String::new(),
1137 keep_issue_open: false,
1138 ignore_checks: false,
1139 bypass_rules: false,
1140 })
1141 .await?;
1142 match merged {
1143 Outcome::Ok(_) => Ok(()),
1144 // Most likely the branch moved in the moment between: let go, and
1145 // the next look at it will catch up and try again.
1146 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
1147 self.db
1148 .prepare(
1149 "UPDATE pulls SET working_on = NULL, working_until = NULL
1150 WHERE id = ? AND working_on = 'merge'",
1151 )
1152 .bind(&[pull.id.as_str().into()])?
1153 .run()
1154 .await?;
1155 Ok(())
1156 }
1157 Outcome::Fail(failure) => {
1158 self.stall(StallArgs {
1159 pull_id: pull.id.clone(),
1160 reason: format!("g1t could not merge this: {}", failure.message),
1161 by: None,
1162 })
1163 .await?;
1164 Ok(())
1165 }
1166 }
1167 }
1168
1169 /// Records that a merge was asked for while the pull request was
1170 /// behind, and announces it so that the runner brings it up to date.
1171 pub(crate) async fn request_landing(
1172 &self,
1173 pull: &Pull,
1174 actor: &User,
1175 keep_issue_open: bool,
1176 ) -> Result<()> {
1177 let now = now_ms();
1178 let request = serde_json::to_string(&LandRequest {
1179 actor: actor.clone(),
1180 keep_issue_open,
1181 })?;
1182 self.db
1183 .prepare(
1184 "UPDATE pulls
1185 SET land_requested = ?, land_requested_at = ?, stalled = NULL,
1186 working_on = 'catch_up', working_until = ?
1187 WHERE id = ?",
1188 )
1189 .bind(&[
1190 request.into(),
1191 rfc3339(now).into(),
1192 rfc3339(now + CATCH_UP_MINUTES * 60 * 1000).into(),
1193 pull.id.as_str().into(),
1194 ])?
1195 .run()
1196 .await?;
1197 self.publish(
1198 "pull.merge_requested",
1199 &pull.repo_id,
1200 actor,
1201 Self::pull_event(pull),
1202 )
1203 .await
1204 }
1205
1206 /// The merge waiting on a pull request, if one was asked for recently
1207 /// enough to still stand.
1208 async fn land_request(&self, pull_id: &str) -> Result<Option<LandRequest>> {
1209 let row = self.land_row(pull_id).await?;
1210 let oldest = rfc3339(now_ms().saturating_sub(CATCH_UP_MINUTES * 60 * 1000));
1211 Ok(row
1212 .filter(|row| {
1213 row.land_requested_at
1214 .as_deref()
1215 .is_some_and(|at| at >= oldest.as_str())
1216 })
1217 .and_then(|row| row.land_requested)
1218 .and_then(|request| serde_json::from_str(&request).ok()))
1219 }
1220
1221 /// Whether a merge is waiting on a pull request, and why g1t stopped
1222 /// working on it if it did.
1223 pub(crate) async fn landing_state(&self, pull_id: &str) -> Result<(bool, Option<String>)> {
1224 let stalled = self.land_row(pull_id).await?.and_then(|row| row.stalled);
1225 Ok((self.land_request(pull_id).await?.is_some(), stalled))
1226 }
1227
1228 /// The merge waiting on a pull request and why g1t stopped, as read
1229 /// for this request (prefetch.rs) or now.
1230 async fn land_row(&self, pull_id: &str) -> Result<Option<LandRow>> {
1231 if let Some(found) = self.prefetched_pull(pull_id) {
1232 return found.first::<LandRow>(Slot::Pull);
1233 }
1234 self.db
1235 .prepare("SELECT land_requested, land_requested_at, stalled FROM pulls WHERE id = ?")
1236 .bind(&[pull_id.into()])?
1237 .first::<LandRow>(None)
1238 .await
1239 }
1240
1241 async fn forget_landing(&self, pull_id: &str) -> Result<()> {
1242 self.db
1243 .prepare(
1244 "UPDATE pulls SET land_requested = NULL, land_requested_at = NULL WHERE id = ?",
1245 )
1246 .bind(&[pull_id.into()])?
1247 .run()
1248 .await?;
1249 Ok(())
1250 }
1251
1252 /// Lands a pull request whose head has just moved, if a merge of it was
1253 /// waiting for exactly that. What it was caught up to was already
1254 /// checked and reviewed apart from the merge, so the checks are not
1255 /// waited for again; a repository that wants them rerun turns on
1256 /// "require up to date", and then nothing is landed this way.
1257 pub(crate) async fn land_if_requested(&self, pull_id: &str) -> Result<()> {
1258 let Some(request) = self.land_request(pull_id).await? else {
1259 return Ok(());
1260 };
1261 self.forget_landing(pull_id).await?;
1262 let Some(pull) = self.pull_by_id(pull_id).await? else {
1263 return Ok(());
1264 };
1265 let repo: Outcome<Repo> = g1t_kit::call(
1266 &self.repos,
1267 "get_by_id",
1268 &GetByIdArgs {
1269 id: pull.repo_id.clone(),
1270 viewer: Some(request.actor.clone()),
1271 },
1272 )
1273 .await?;
1274 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
1275 return Ok(());
1276 };
1277 let merged = self
1278 .merge_pull(PullActionArgs {
1279 actor: request.actor,
1280 repo: RepoPath {
1281 namespace: repo.namespace,
1282 name: repo.name,
1283 },
1284 number: pull.number,
1285 summary: String::new(),
1286 keep_issue_open: request.keep_issue_open,
1287 ignore_checks: true,
1288 bypass_rules: false,
1289 })
1290 .await?;
1291 if let Outcome::Fail(failure) = merged {
1292 self.stall(StallArgs {
1293 pull_id: pull.id,
1294 reason: format!(
1295 "It was brought up to date but could not be merged: {}",
1296 failure.message
1297 ),
1298 by: None,
1299 })
1300 .await?;
1301 }
1302 Ok(())
1303 }
1304
1305 /// What the runner needs to bring a pull request up to date for a merge
1306 /// that is waiting on it.
1307 pub(crate) async fn catch_up_job(&self, a: CatchUpJobArgs) -> Result<Option<LifecycleJob>> {
1308 if self.land_request(&a.pull_id).await?.is_none() {
1309 return Ok(None);
1310 }
1311 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
1312 return Ok(None);
1313 };
1314 // Its owner (whoever asked g1t for it, or its author) can read both
1315 // the repository and the pull request's source.
1316 let repo: Outcome<Repo> = g1t_kit::call(
1317 &self.repos,
1318 "get_by_id",
1319 &GetByIdArgs {
1320 id: pull.repo_id.clone(),
1321 viewer: self.owner_viewer(&pull).await?,
1322 },
1323 )
1324 .await?;
1325 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
1326 return Ok(None);
1327 };
1328 let path = RepoPath {
1329 namespace: repo.namespace,
1330 name: repo.name,
1331 };
1332 let issue = match pull.issue {
1333 Some(number) => self.issue(&pull.repo_id, number).await?,
1334 None => None,
1335 };
1336 let base = pull.base_branch(&repo.default_branch).to_owned();
1337 let feedback = self.conflict_note(&pull, &base).await?;
1338 Ok(Some(LifecycleJob {
1339 pull_id: pull.id,
1340 source: pull.fork.unwrap_or_else(|| path.clone()),
1341 repo: path,
1342 number: pull.number,
1343 author: pull.requested_by.unwrap_or(pull.author),
1344 branch: pull.branch,
1345 default_branch: base,
1346 title: pull.title,
1347 description: pull.body.unwrap_or_default(),
1348 issue,
1349 feedback,
1350 round: 0,
1351 }))
1352 }
1353
1354 /// For an agent catching up: the files g1t already knows conflict, so
1355 /// it reads them first. Empty when none are known.
1356 pub(crate) async fn conflict_note(&self, pull: &Pull, default_branch: &str) -> Result<String> {
1357 Ok(match self.conflicting_files(pull).await? {
1358 Some(files) if !files.is_empty() => format!(
1359 "g1t found ahead of time that merging {default_branch} into this pull request conflicts in these files: {}.",
1360 files.join(", ")
1361 ),
1362 _ => String::new(),
1363 })
1364 }
1365
1366 /// Stops seeing a pull request through until a person steps in. The
1367 /// first stop is published (`pull.stalled`), which tells its people
1368 /// that it needs them; a stop on one already stopped only says why.
1369 pub(crate) async fn stall(&self, a: StallArgs) -> Result<bool> {
1370 let before = self.pull_by_id(&a.pull_id).await?;
1371 let was_stalled = self.is_stalled(&a.pull_id).await?;
1372 self.db
1373 .prepare(
1374 "UPDATE pulls
1375 SET stalled = ?, working_on = NULL, working_until = NULL,
1376 land_requested = NULL, land_requested_at = NULL
1377 WHERE id = ? AND status = 'open'",
1378 )
1379 .bind(&[a.reason.trim().into(), a.pull_id.as_str().into()])?
1380 .run()
1381 .await?;
1382 self.db
1383 .prepare("UPDATE pulls SET stage = 'needs_you', stage_detail = ? WHERE id = ?")
1384 .bind(&[a.reason.trim().into(), a.pull_id.as_str().into()])?
1385 .run()
1386 .await?;
1387 if let Some(pull) = before.filter(|pull| !was_stalled && pull.status == PullStatus::Open) {
1388 self.publish_as(
1389 "pull.stalled",
1390 &pull.repo_id,
1391 a.by.clone(),
1392 g1t_contracts::events::PullEvent {
1393 detail: Some(a.reason.trim().to_owned()),
1394 ..Self::pull_event(&pull)
1395 },
1396 )
1397 .await?;
1398 }
1399 Ok(true)
1400 }
1401
1402 /// Whether g1t has stopped seeing the pull request through.
1403 pub(crate) async fn is_stalled(&self, pull_id: &str) -> Result<bool> {
1404 Ok(self
1405 .db
1406 .prepare("SELECT 1 AS value FROM pulls WHERE id = ? AND stalled IS NOT NULL")
1407 .bind(&[pull_id.into()])?
1408 .first::<u32>(Some("value"))
1409 .await?
1410 .is_some())
1411 }
1412
1413 /// Says that a pull request g1t had stopped on is going again, which
1414 /// closes what it was waiting on a person for.
1415 pub(crate) async fn announce_resumed(&self, pull_id: &str, actor: Option<String>) -> Result<()> {
1416 if let Some(pull) = self.pull_by_id(pull_id).await? {
1417 self.publish_as("pull.resumed", &pull.repo_id, actor, Self::pull_event(&pull)).await?;
1418 }
1419 Ok(())
1420 }
1421
1422 pub(crate) async fn managed_pulls(&self, a: ManagedPullsArgs) -> Result<Vec<String>> {
1423 let rows = self
1424 .db
1425 .prepare(
1426 "SELECT id AS value FROM pulls
1427 WHERE status = 'open' AND managed = 1 AND stalled IS NULL
1428 AND (?1 IS NULL OR repo_id = ?1)
1429 ORDER BY updated_at DESC LIMIT ?2",
1430 )
1431 .bind(&[
1432 a.repo_id.map_or(JsValue::NULL, JsValue::from),
1433 MANAGED_PAGE.into(),
1434 ])?
1435 .all()
1436 .await?
1437 .results::<ValueRow>()?;
1438 Ok(rows.into_iter().map(|row| row.value).collect())
1439 }
1440}
1441
1442#[cfg(test)]
1443mod tests {
1444 use super::*;
1445
1446 const MAX_REVISIONS: u32 = 2;
1447
1448 #[test]
1449 fn a_verdict_from_whoever_asked_for_g1t_s_change_is_not_someone_else_s() {
1450 use crate::rows::stored::{ASKER, G1T, pull};
1451 let made = pull(G1T, Some(ASKER));
1452 assert!(made_by_g1t(&made));
1453 // The person it was made for is held to what an author was: their
1454 // request for changes does not send g1t back as a reviewer's
1455 // would, and their approval does not lift a hold.
1456 assert!(!from_someone_else(&made, ASKER.0));
1457 assert!(!from_someone_else(&made, AGENT_ID));
1458 assert!(from_someone_else(&made, "usr_reviewer"));
1459 // Anyone's own pull request, the same.
1460 let own = pull(ASKER, None);
1461 assert!(!from_someone_else(&own, ASKER.0));
1462 assert!(from_someone_else(&own, "usr_reviewer"));
1463 }
1464
1465 #[test]
1466 fn approvals_the_repository_wants_are_waited_for() {
1467 let short = || Facts {
1468 review: reviewed(Some(Verdict::Approve)),
1469 approvals_missing: Some("This repository requires 1 approving review.".to_owned()),
1470 ..facts()
1471 };
1472 let (lifecycle, next) = decide(short());
1473 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1474 assert_eq!(
1475 lifecycle.detail,
1476 "This repository requires 1 approving review."
1477 );
1478 assert!(matches!(next, Next::Wait));
1479 // Not even a repository that merges by itself merges without them.
1480 let automatic = Facts {
1481 auto_merge: true,
1482 ..short()
1483 };
1484 assert_eq!(outcome(automatic), (Stage::NeedsYou, "wait"));
1485 }
1486
1487 #[test]
1488 fn a_repository_can_leave_review_to_people() {
1489 let unreviewed = Facts {
1490 agent_review: false,
1491 ..facts()
1492 };
1493 assert_eq!(outcome(unreviewed), (Stage::Ready, "wait"));
1494 let failing = Facts {
1495 agent_review: false,
1496 check_status: Some(CheckStatus::Failed),
1497 ..facts()
1498 };
1499 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1500 }
1501
1502 #[test]
1503 fn a_repository_sets_how_often_the_author_is_sent_back() {
1504 let never = Facts {
1505 max_revisions: 0,
1506 check_status: Some(CheckStatus::Failed),
1507 ..facts()
1508 };
1509 assert_eq!(outcome(never), (Stage::NeedsYou, "wait"));
1510 }
1511
1512 /// Checks on a head commit: `(context, state)` statuses, against the
1513 /// required check names.
1514 fn checks(statuses: &[(&str, &str)], required: &[&str]) -> WorkflowFacts {
1515 let statuses: Vec<CommitStatus> = statuses
1516 .iter()
1517 .map(|(context, state)| CommitStatus {
1518 context: (*context).to_owned(),
1519 state: (*state).to_owned(),
1520 description: None,
1521 target_url: None,
1522 updated_at: String::new(),
1523 source: None,
1524 })
1525 .collect();
1526 let required: Vec<String> = required.iter().map(|name| (*name).to_owned()).collect();
1527 WorkflowFacts::of(&statuses, &required)
1528 }
1529
1530 /// A pull request that is ready for review, whose required check
1531 /// passed, and nothing else yet.
1532 fn facts() -> Facts {
1533 Facts {
1534 draft: false,
1535 stalled: None,
1536 working_on: None,
1537 check_status: None,
1538 review_pending: false,
1539 revisions: 0,
1540 review: None,
1541 behind: false,
1542 conflicting: false,
1543 auto_merge: false,
1544 require_up_to_date: false,
1545 agent_review: true,
1546 max_revisions: MAX_REVISIONS,
1547 approvals_missing: None,
1548 person_request: None,
1549 queued: None,
1550 workflows: checks(&[("CI / pull_request", "success")], &["CI"]),
1551 low_confidence: None,
1552 }
1553 }
1554
1555 #[test]
1556 fn failed_checks_send_the_agent_back_and_running_ones_wait() {
1557 let failed = Facts {
1558 workflows: checks(&[("CI / pull_request", "failure")], &["CI"]),
1559 ..facts()
1560 };
1561 let (lifecycle, next) = decide(failed);
1562 assert!(matches!(next, Next::Revise(Feedback::FailedWorkflows)));
1563 assert!(lifecycle.detail.contains("CI / pull_request failed"));
1564 let running = Facts {
1565 workflows: checks(&[("CI / pull_request", "pending")], &["CI"]),
1566 ..facts()
1567 };
1568 let (lifecycle, next) = decide(running);
1569 assert!(matches!(next, Next::Wait));
1570 assert!(lifecycle.detail.contains("Waiting for CI / pull_request"));
1571 }
1572
1573 #[test]
1574 fn a_check_the_branch_does_not_require_is_fixed_but_does_not_hold_it_for_ever() {
1575 // Lint is not required: the agent is still sent back to fix it...
1576 let lint = checks(&[("CI / pull_request", "success"), ("Lint / pull_request", "failure")], &["CI"]);
1577 let failing = Facts { workflows: lint.clone(), ..facts() };
1578 assert_eq!(outcome(failing), (Stage::Revising, "revise for workflows"));
1579 // ...but once it is out of revisions, Lint no longer holds it.
1580 let exhausted = Facts {
1581 workflows: lint,
1582 revisions: MAX_REVISIONS,
1583 review: reviewed(Some(Verdict::Approve)),
1584 ..facts()
1585 };
1586 assert_eq!(outcome(exhausted), (Stage::Ready, "wait"));
1587 // A required check that still fails asks a person.
1588 let required = Facts {
1589 workflows: checks(&[("CI / pull_request", "failure")], &["CI"]),
1590 revisions: MAX_REVISIONS,
1591 ..facts()
1592 };
1593 let (lifecycle, next) = decide(required);
1594 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1595 assert_eq!(lifecycle.detail, "The required check CI still fails after the agent revised twice.");
1596 assert!(matches!(next, Next::Wait));
1597 }
1598
1599 #[test]
1600 fn a_required_check_that_has_not_reported_is_waited_for() {
1601 let missing = Facts {
1602 workflows: checks(&[("CI / pull_request", "success")], &["CI", "Deploy"]),
1603 review: reviewed(Some(Verdict::Approve)),
1604 auto_merge: true,
1605 ..facts()
1606 };
1607 let (lifecycle, next) = decide(missing);
1608 assert_eq!(lifecycle.stage, Stage::Checking);
1609 assert_eq!(lifecycle.detail, "Waiting for the required check Deploy to report on its latest commit.");
1610 assert!(matches!(next, Next::Wait), "auto-merge must not land it");
1611 }
1612
1613 #[test]
1614 fn a_queued_pull_request_waits_in_the_queue() {
1615 let queued = Facts {
1616 review: reviewed(Some(Verdict::Approve)),
1617 queued: Some((QueueState::Testing, vec![12, 14])),
1618 auto_merge: true,
1619 ..facts()
1620 };
1621 let (lifecycle, next) = decide(queued);
1622 assert_eq!(lifecycle.stage, Stage::Queued);
1623 assert!(lifecycle.detail.contains("#12, #14"));
1624 assert!(matches!(next, Next::Wait));
1625 }
1626
1627 fn asked_by_a_person() -> Option<PersonRequest> {
1628 Some(PersonRequest {
1629 author_id: "usr_reviewer".to_owned(),
1630 author_name: "g1t-reviewer".to_owned(),
1631 created_at: "2026-10-02T11:00:00.000Z".to_owned(),
1632 })
1633 }
1634
1635 #[test]
1636 fn a_person_asking_for_changes_sends_the_agent_back() {
1637 let asked = Facts {
1638 review: reviewed(Some(Verdict::Approve)),
1639 approvals_missing: Some("A reviewer has asked for changes.".to_owned()),
1640 person_request: asked_by_a_person(),
1641 ..facts()
1642 };
1643 let (lifecycle, _) = decide(Facts {
1644 person_request: asked_by_a_person(),
1645 ..facts()
1646 });
1647 assert!(lifecycle.detail.starts_with("g1t-reviewer asked for changes"));
1648 assert_eq!(outcome(asked), (Stage::Revising, "revise for a person"));
1649 }
1650
1651 #[test]
1652 fn a_person_is_asked_once_the_revisions_run_out() {
1653 let exhausted = Facts {
1654 person_request: asked_by_a_person(),
1655 revisions: MAX_REVISIONS,
1656 ..facts()
1657 };
1658 assert_eq!(outcome(exhausted), (Stage::NeedsYou, "wait"));
1659 }
1660
1661 fn reviewed(verdict: Option<Verdict>) -> Option<FinishedReview> {
1662 Some(FinishedReview {
1663 finished_at: "2026-10-02T10:00:00.000Z".to_owned(),
1664 verdict,
1665 })
1666 }
1667
1668 /// The stage, and a word for the step to take.
1669 fn outcome(facts: Facts) -> (Stage, &'static str) {
1670 let (lifecycle, next) = decide(facts);
1671 let step = match next {
1672 Next::Wait => "wait",
1673 Next::Review => "review",
1674 Next::Revise(Feedback::FailedChecks) => "revise for checks",
1675 Next::Revise(Feedback::FailedWorkflows) => "revise for workflows",
1676 Next::Revise(Feedback::Review(_)) => "revise for review",
1677 Next::Revise(Feedback::Person(_)) => "revise for a person",
1678 Next::CatchUp => "catch up",
1679 Next::Merge => "merge",
1680 };
1681 (lifecycle.stage, step)
1682 }
1683
1684 #[test]
1685 fn nothing_is_started_while_the_agent_is_still_working() {
1686 let draft = Facts {
1687 draft: true,
1688 check_status: None,
1689 ..facts()
1690 };
1691 assert_eq!(outcome(draft), (Stage::Working, "wait"));
1692 }
1693
1694 #[test]
1695 fn checks_come_before_review() {
1696 let unchecked = Facts {
1697 workflows: checks(&[], &["CI"]),
1698 ..facts()
1699 };
1700 assert_eq!(outcome(unchecked), (Stage::Checking, "wait"));
1701 let running = Facts {
1702 workflows: checks(&[("CI / pull_request", "pending")], &["CI"]),
1703 ..facts()
1704 };
1705 assert_eq!(outcome(running), (Stage::Checking, "wait"));
1706 assert_eq!(outcome(facts()), (Stage::Reviewing, "review"));
1707 }
1708
1709 #[test]
1710 fn a_branch_that_requires_no_checks_goes_straight_to_review() {
1711 let unchecked = Facts {
1712 workflows: WorkflowFacts::default(),
1713 ..facts()
1714 };
1715 assert_eq!(outcome(unchecked), (Stage::Reviewing, "review"));
1716 }
1717
1718 #[test]
1719 fn failing_in_the_merge_queue_sends_the_author_back() {
1720 let failed = Facts {
1721 check_status: Some(CheckStatus::Failed),
1722 ..facts()
1723 };
1724 let (lifecycle, next) = decide(failed);
1725 assert_eq!(lifecycle.detail, "It failed in the merge queue. The agent is being sent back to fix it.");
1726 assert!(matches!(next, Next::Revise(Feedback::FailedChecks)));
1727 }
1728
1729 #[test]
1730 fn checks_that_could_not_run_are_a_persons_problem() {
1731 let errored = Facts {
1732 check_status: Some(CheckStatus::Errored),
1733 ..facts()
1734 };
1735 assert_eq!(outcome(errored), (Stage::NeedsYou, "wait"));
1736 }
1737
1738 #[test]
1739 fn a_review_asking_for_changes_sends_the_author_back() {
1740 let changes = Facts {
1741 review: reviewed(Some(Verdict::RequestChanges)),
1742 ..facts()
1743 };
1744 assert_eq!(outcome(changes), (Stage::Revising, "revise for review"));
1745 }
1746
1747 #[test]
1748 fn the_author_is_sent_back_only_so_many_times() {
1749 let failing = Facts {
1750 check_status: Some(CheckStatus::Failed),
1751 revisions: MAX_REVISIONS,
1752 ..facts()
1753 };
1754 assert_eq!(outcome(failing), (Stage::NeedsYou, "wait"));
1755 let unconvinced = Facts {
1756 review: reviewed(Some(Verdict::RequestChanges)),
1757 revisions: MAX_REVISIONS,
1758 ..facts()
1759 };
1760 assert_eq!(outcome(unconvinced), (Stage::NeedsYou, "wait"));
1761 // One short of the limit still gets another go.
1762 let once = Facts {
1763 check_status: Some(CheckStatus::Failed),
1764 revisions: MAX_REVISIONS - 1,
1765 ..facts()
1766 };
1767 assert_eq!(outcome(once), (Stage::Revising, "revise for checks"));
1768 }
1769
1770 #[test]
1771 fn a_review_that_could_not_be_written_is_not_retried() {
1772 let broken = Facts {
1773 review: reviewed(None),
1774 ..facts()
1775 };
1776 assert_eq!(outcome(broken), (Stage::NeedsYou, "wait"));
1777 }
1778
1779 #[test]
1780 fn being_behind_only_holds_a_change_up_where_the_repository_says_so() {
1781 let behind = || Facts {
1782 review: reviewed(Some(Verdict::Approve)),
1783 behind: true,
1784 ..facts()
1785 };
1786 // By default it is ready as it is; merging brings it up to date.
1787 assert_eq!(outcome(behind()), (Stage::Ready, "wait"));
1788 let strict = Facts {
1789 require_up_to_date: true,
1790 ..behind()
1791 };
1792 assert_eq!(outcome(strict), (Stage::CatchingUp, "catch up"));
1793 let current = Facts {
1794 review: reviewed(Some(Verdict::Approve)),
1795 ..facts()
1796 };
1797 assert_eq!(outcome(current), (Stage::Ready, "wait"));
1798 }
1799
1800 #[test]
1801 fn a_ready_change_lands_by_itself_only_where_the_repository_says_so() {
1802 let ready = || Facts {
1803 review: reviewed(Some(Verdict::Approve)),
1804 ..facts()
1805 };
1806 assert_eq!(outcome(ready()), (Stage::Ready, "wait"));
1807 let automatic = Facts {
1808 auto_merge: true,
1809 ..ready()
1810 };
1811 assert_eq!(outcome(automatic), (Stage::Ready, "merge"));
1812 // One that is behind is merged too: merging brings it up to date.
1813 let behind = Facts {
1814 auto_merge: true,
1815 behind: true,
1816 ..ready()
1817 };
1818 assert_eq!(outcome(behind), (Stage::Ready, "merge"));
1819 // Unless the repository wants it caught up and checked again first.
1820 let strict = Facts {
1821 auto_merge: true,
1822 behind: true,
1823 require_up_to_date: true,
1824 ..ready()
1825 };
1826 assert_eq!(outcome(strict), (Stage::CatchingUp, "catch up"));
1827 // Nothing short of approved is merged, whatever the setting.
1828 let failing = Facts {
1829 auto_merge: true,
1830 check_status: Some(CheckStatus::Failed),
1831 ..ready()
1832 };
1833 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1834 let unreviewed = Facts {
1835 auto_merge: true,
1836 ..facts()
1837 };
1838 assert_eq!(outcome(unreviewed), (Stage::Reviewing, "review"));
1839 }
1840
1841 #[test]
1842 fn a_conflict_found_ahead_of_time_is_resolved_before_merging() {
1843 let conflicting = || Facts {
1844 review: reviewed(Some(Verdict::Approve)),
1845 behind: true,
1846 conflicting: true,
1847 ..facts()
1848 };
1849 // Even where the repository would merge one that is merely behind.
1850 assert_eq!(outcome(conflicting()), (Stage::CatchingUp, "catch up"));
1851 let automatic = Facts {
1852 auto_merge: true,
1853 ..conflicting()
1854 };
1855 assert_eq!(outcome(automatic), (Stage::CatchingUp, "catch up"));
1856 // Failed checks come first: a revision merges the branch in too.
1857 let failing = Facts {
1858 check_status: Some(CheckStatus::Failed),
1859 ..conflicting()
1860 };
1861 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1862 }
1863
1864 #[test]
1865 fn catching_up_reruns_the_checks_but_not_the_review() {
1866 // The merge moved the head, so the checks are waited for again.
1867 let merged_in = Facts {
1868 review: reviewed(Some(Verdict::Approve)),
1869 workflows: checks(&[], &["CI"]),
1870 ..facts()
1871 };
1872 assert_eq!(outcome(merged_in), (Stage::Checking, "wait"));
1873 }
1874
1875 #[test]
1876 fn a_step_under_way_is_not_started_again() {
1877 for (step, stage) in [
1878 ("review", Stage::Reviewing),
1879 ("revision", Stage::Revising),
1880 ("catch_up", Stage::CatchingUp),
1881 ] {
1882 let busy = Facts {
1883 working_on: Some(step.to_owned()),
1884 // Whatever else is true, the step in hand comes first.
1885 check_status: Some(CheckStatus::Failed),
1886 ..facts()
1887 };
1888 assert_eq!(outcome(busy), (stage, "wait"));
1889 }
1890 let asked = Facts {
1891 review_pending: true,
1892 ..facts()
1893 };
1894 assert_eq!(outcome(asked), (Stage::Reviewing, "wait"));
1895 }
1896
1897 #[test]
1898 fn a_low_confidence_change_waits_for_a_person_instead_of_merging() {
1899 let held = || Facts {
1900 review: reviewed(Some(Verdict::Approve)),
1901 auto_merge: true,
1902 low_confidence: Some("tests not added, 3 revisions".to_owned()),
1903 ..facts()
1904 };
1905 let (lifecycle, next) = decide(held());
1906 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1907 assert!(matches!(next, Next::Wait), "auto-merge must not land it");
1908 assert_eq!(
1909 lifecycle.detail,
1910 "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."
1911 );
1912 // Without auto-merge it needs someone too, and says why.
1913 assert_eq!(outcome(Facts { auto_merge: false, ..held() }), (Stage::NeedsYou, "wait"));
1914 // Not held (the setting is off, or a person approved): it lands.
1915 assert_eq!(outcome(Facts { low_confidence: None, ..held() }), (Stage::Ready, "merge"));
1916 // It holds only a change that is otherwise ready: what comes first,
1917 // such as failed checks, is still dealt with first.
1918 let failing = Facts {
1919 check_status: Some(CheckStatus::Failed),
1920 ..held()
1921 };
1922 assert_eq!(outcome(failing), (Stage::Revising, "revise for checks"));
1923 let unapproved = Facts {
1924 approvals_missing: Some("This repository requires 1 approving review.".to_owned()),
1925 ..held()
1926 };
1927 assert_eq!(decide(unapproved).0.detail, "This repository requires 1 approving review.");
1928 }
1929
1930 #[test]
1931 fn once_stopped_it_stays_stopped() {
1932 let (lifecycle, next) = decide(Facts {
1933 stalled: Some("The agent could not catch up.".to_owned()),
1934 review: reviewed(Some(Verdict::Approve)),
1935 behind: true,
1936 ..facts()
1937 });
1938 assert_eq!(lifecycle.stage, Stage::NeedsYou);
1939 assert_eq!(lifecycle.detail, "The agent could not catch up.");
1940 assert!(matches!(next, Next::Wait));
1941 }
1942}