pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/queue.rs

770 lines29,877 bytesCodeBlame

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

Agents as a team: lifecycle, merge queue, billing and a new shell1//! The merge queue: pull requests are tested together with the ones ahead
2//! of them, and only a combination that passed reaches the default branch.
3//!
4//! A batch of up to [`BATCH`] entries is tested at once, speculatively: one
5//! sandbox per entry builds the default branch with that entry and every
6//! entry ahead of it merged in, runs all of their acceptance checks, and
7//! pushes the result to `g1t-queue/<entry>`. Entries land in order, each by
8//! moving the default branch to its tested state, once everything ahead has
9//! landed. One that fails leaves the queue with a failed check run, which
10//! sends a g1t agent back to fix it; the entries behind it are tested again
11//! without it.
12
13use futures_util::future::try_join_all;
14use g1t_contracts::events::{ChecksEvent, QueueChanged};
15use g1t_contracts::repos::{GetByIdArgs, HeadArgs, LandArgs, Landed, Repo, RepoPath};
16use g1t_contracts::time::rfc3339;
17use g1t_contracts::work::*;
18use g1t_contracts::{FailureCode, Outcome, User, new_id};
19use g1t_kit::now_ms;
20use serde::Deserialize;
21use worker::Result;
22use worker::wasm_bindgen::JsValue;
23
24use crate::Work;
25use crate::checks::{hash, new_token};
26
27/// How many entries are tested at once.
28const BATCH: usize = 4;
29/// How long a batch may take before it is tested again.
30const TESTING_MINUTES: u64 = 45;
31/// How many entries that left the queue are shown.
32const RECENT: u32 = 20;
33/// How many completed issues' checks make up the default branch's contract.
34const CONTRACT_ISSUES: u32 = 30;
35/// Where tested states are pushed, in the repository itself.
36const BRANCH_PREFIX: &str = "g1t-queue/";
37
38#[derive(Clone, Deserialize)]
Sidebar: the panels really slide39pub(crate) struct EntryRow {
Agents as a team: lifecycle, merge queue, billing and a new shell40 id: String,
41 repo_id: String,
42 pull_id: String,
43 number: u32,
44 state: String,
45 enqueued_by: String,
46 keep_issue_open: u32,
47 head_commit: Option<String>,
48 base_commit: Option<String>,
49 ahead: Option<String>,
50 combined_commit: Option<String>,
51 results: Option<String>,
52 error: Option<String>,
53 tested_at: Option<String>,
54 created_at: String,
55 finished_at: Option<String>,
56}
57
58impl EntryRow {
59 fn state(&self) -> QueueState {
60 match self.state.as_str() {
61 "testing" => QueueState::Testing,
62 "passed" => QueueState::Passed,
63 "failed" => QueueState::Failed,
64 "landed" => QueueState::Landed,
65 "removed" => QueueState::Removed,
66 _ => QueueState::Waiting,
67 }
68 }
69
70 fn actor(&self) -> Option<User> {
71 serde_json::from_str(&self.enqueued_by).ok()
72 }
73
74 fn ahead(&self) -> Vec<u32> {
75 self.ahead
76 .as_deref()
77 .and_then(|ahead| serde_json::from_str(ahead).ok())
78 .unwrap_or_default()
79 }
80
81 fn branch(&self) -> String {
82 format!("{BRANCH_PREFIX}{}", self.id)
83 }
84}
85
86fn minutes_ago(minutes: u64) -> String {
87 rfc3339(now_ms().saturating_sub(minutes * 60 * 1000))
88}
89
90impl Work {
91 async fn entries(&self, repo_id: &str, active: bool) -> Result<Vec<EntryRow>> {
92 let sql = if active {
93 "SELECT * FROM queue_entries
94 WHERE repo_id = ? AND state IN ('waiting', 'testing', 'passed')
95 ORDER BY created_at, id"
96 } else {
97 "SELECT * FROM queue_entries
98 WHERE repo_id = ? AND state IN ('failed', 'landed', 'removed')
99 ORDER BY finished_at DESC LIMIT ?"
100 };
101 let statement = self.db.prepare(sql);
102 let statement = if active {
103 statement.bind(&[repo_id.into()])?
104 } else {
105 statement.bind(&[repo_id.into(), RECENT.into()])?
106 };
107 statement.all().await?.results::<EntryRow>()
108 }
109
110 /// The entry a pull request has in the queue now, if any.
111 pub(crate) async fn queued_entry(&self, pull_id: &str) -> Result<Option<(QueueState, Vec<u32>)>> {
112 let row = self
113 .db
114 .prepare(
115 "SELECT * FROM queue_entries
116 WHERE pull_id = ? AND state IN ('waiting', 'testing', 'passed') LIMIT 1",
117 )
118 .bind(&[pull_id.into()])?
119 .first::<EntryRow>(None)
120 .await?;
121 Ok(row.map(|row| (row.state(), row.ahead())))
122 }
123
124 async fn changed(&self, repo_id: &str) -> Result<()> {
125 self.publish_as(
126 "queue.changed",
127 repo_id,
128 None,
129 QueueChanged {
130 repo_id: repo_id.to_owned(),
131 },
132 )
133 .await
134 }
135
136 /// Puts a pull request that may merge into the queue instead. Merging
137 /// it again while it is queued changes nothing.
138 pub(crate) async fn enqueue(
139 &self,
140 repo: &Repo,
141 pull: &Pull,
142 actor: &User,
143 keep_issue_open: bool,
144 ) -> Result<Outcome<Pull>> {
145 if self.queued_entry(&pull.id).await?.is_some() {
146 return Ok(Outcome::Ok(pull.clone()));
147 }
148 let now = now_ms();
149 self.db
150 .prepare(
151 "INSERT INTO queue_entries
152 (id, repo_id, pull_id, number, state, enqueued_by, keep_issue_open, created_at)
153 VALUES (?, ?, ?, ?, 'waiting', ?, ?, ?)",
154 )
155 .bind(&[
156 new_id("qen", now).into(),
157 repo.id.as_str().into(),
158 pull.id.as_str().into(),
159 pull.number.into(),
160 serde_json::to_string(actor)?.into(),
161 u32::from(keep_issue_open).into(),
162 rfc3339(now).into(),
163 ])?
164 .run()
165 .await?;
166 let who = (actor.id.as_str(), actor.username.as_str());
167 self.note(&repo.id, pull.number, who, "added this to the merge queue")
168 .await?;
169 self.changed(&repo.id).await?;
170 Ok(Outcome::Ok(pull.clone()))
171 }
172
173 pub(crate) async fn queue(&self, a: QueueArgs) -> Result<Outcome<QueueView>> {
174 let repo = match self.repo(&a.repo, &a.viewer).await? {
175 Outcome::Ok(repo) => repo,
176 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
177 };
178 let enabled = self.settings(&repo.id).await?.merge_queue;
179 let (active, recent) = (self.entries(&repo.id, true).await?, self.entries(&repo.id, false).await?);
180 let numbers: Vec<u32> = active.iter().chain(&recent).map(|row| row.number).collect();
181 let pulls = try_join_all(numbers.iter().map(|number| self.pull(&repo.id, *number))).await?;
182 let view = |rows: Vec<EntryRow>, pulls: &[Option<Pull>]| -> Vec<QueueEntry> {
183 rows.into_iter()
184 .zip(pulls)
185 .map(|(row, pull)| QueueEntry {
186 state: row.state(),
187 ahead: row.ahead(),
188 enqueued_by: row.actor().map_or_else(|| "g1t".to_owned(), |actor| actor.username),
189 title: pull.as_ref().map(|p| p.title.clone()).unwrap_or_default(),
190 agent: pull.as_ref().map(|p| p.agent.clone()).unwrap_or_default(),
191 results: row
192 .results
193 .as_deref()
194 .and_then(|results| serde_json::from_str(results).ok())
195 .unwrap_or_default(),
196 id: row.id,
197 number: row.number,
198 base_commit: row.base_commit,
199 combined_commit: row.combined_commit,
200 error: row.error,
201 created_at: row.created_at,
202 finished_at: row.finished_at,
203 })
204 .collect()
205 };
206 let split = active.len();
207 Ok(Outcome::Ok(QueueView {
208 enabled,
209 active: view(active, &pulls[..split]),
210 recent: view(recent, &pulls[split..]),
211 }))
212 }
213
214 /// Takes a pull request out of the queue, by a member's hand.
215 pub(crate) async fn remove_from_queue(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
216 let viewer = Some(a.actor.clone());
217 let (repo, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
218 Outcome::Ok(found) => found,
219 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
220 };
221 if !a.actor.verified || !a.actor.is_member(&repo.namespace) {
222 return Ok(Outcome::fail(
223 FailureCode::Forbidden,
224 "Only members of the repository's workspace can change its merge queue.",
225 ));
226 }
227 if self.leave(&repo.id, &pull, QueueState::Removed, None).await? {
228 let who = (a.actor.id.as_str(), a.actor.username.as_str());
229 self.note(&repo.id, pull.number, who, "removed this from the merge queue")
230 .await?;
231 self.changed(&repo.id).await?;
232 }
233 Ok(Outcome::Ok(pull))
234 }
235
236 /// Takes a pull request's entry out of the queue, and sends every entry
237 /// whose tested state included it back to waiting. Whether it had one.
238 pub(crate) async fn leave(
239 &self,
240 repo_id: &str,
241 pull: &Pull,
242 state: QueueState,
243 error: Option<&str>,
244 ) -> Result<bool> {
245 let active = self.entries(repo_id, true).await?;
246 let Some(entry) = active.iter().find(|row| row.pull_id == pull.id) else {
247 return Ok(false);
248 };
249 let now = rfc3339(now_ms());
250 self.db
251 .prepare(
252 "UPDATE queue_entries SET state = ?, error = COALESCE(?, error), finished_at = ?
253 WHERE id = ?",
254 )
255 .bind(&[
256 state.as_str().into(),
257 error.map_or(JsValue::NULL, JsValue::from),
258 now.as_str().into(),
259 entry.id.as_str().into(),
260 ])?
261 .run()
262 .await?;
263 let behind: Vec<&EntryRow> = active
264 .iter()
265 .filter(|row| row.state() != QueueState::Waiting && row.ahead().contains(&pull.number))
266 .collect();
267 self.retest(&behind).await?;
268 Ok(true)
269 }
270
271 /// Sends entries back to waiting, to be tested again.
272 async fn retest(&self, rows: &[&EntryRow]) -> Result<()> {
273 for row in rows {
274 self.db
275 .prepare(
276 "UPDATE queue_entries
277 SET state = 'waiting', token_hash = NULL, combined_commit = NULL,
278 results = NULL, error = NULL, ahead = NULL
279 WHERE id = ? AND state IN ('testing', 'passed')",
280 )
281 .bind(&[row.id.as_str().into()])?
282 .run()
283 .await?;
284 }
285 Ok(())
286 }
287
288 /// The next batch to test, if nothing is being tested now: one job per
289 /// entry, each building the default branch with that entry and every
290 /// entry ahead of it.
291 pub(crate) async fn queue_build(&self, a: QueueBuildArgs) -> Result<Vec<QueueJob>> {
292 let active = self.entries(&a.repo_id, true).await?;
293 // A batch that has taken too long is tested again.
294 let stale = minutes_ago(TESTING_MINUTES);
295 let stuck: Vec<&EntryRow> = active
296 .iter()
297 .filter(|row| {
298 row.state() == QueueState::Testing
299 && row.tested_at.as_deref().is_none_or(|at| at < stale.as_str())
300 })
301 .collect();
302 if !stuck.is_empty() {
303 self.retest(&stuck).await?;
304 return Box::pin(self.queue_build(a)).await;
305 }
306 if active.iter().any(|row| row.state() != QueueState::Waiting) {
307 return Ok(Vec::new());
308 }
309 let batch: Vec<EntryRow> = active.into_iter().take(BATCH).collect();
310 let Some(first) = batch.first() else {
311 return Ok(Vec::new());
312 };
313 let Some(actor) = first.actor() else {
314 return Ok(Vec::new());
315 };
316 let repo: Outcome<Repo> = g1t_kit::call(
317 &self.repos,
318 "get_by_id",
319 &GetByIdArgs {
320 id: a.repo_id.clone(),
321 viewer: Some(actor),
322 },
323 )
324 .await?;
325 let Outcome::Ok(repo) = repo else {
326 return Ok(Vec::new());
327 };
328 let base: Option<String> = g1t_kit::call(
329 &self.repos,
330 "head",
331 &HeadArgs {
332 repo_id: repo.id.clone(),
333 branch: repo.default_branch.clone(),
334 },
335 )
336 .await?;
337 let Some(base) = base else {
338 return Ok(Vec::new());
339 };
340
341 // Each entry's change, as it is now, and the checks it brings.
342 let mut items: Vec<(EntryRow, Pull, QueueStackItem, Vec<String>)> = Vec::new();
343 for row in batch {
344 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
345 continue;
346 };
347 if pull.status != PullStatus::Open {
348 self.leave(&repo.id, &pull, QueueState::Removed, Some("It was closed."))
349 .await?;
350 continue;
351 }
352 let head: Option<String> = g1t_kit::call(
353 &self.repos,
354 "head",
355 &HeadArgs {
356 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
357 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
358 },
359 )
360 .await?;
361 let Some(commit) = head else {
362 continue;
363 };
364 let checks = match pull.issue {
365 Some(number) => self
366 .issue(&repo.id, number)
367 .await?
368 .map(|issue| issue.checks)
369 .unwrap_or_default(),
370 None => Vec::new(),
371 };
372 let source = pull.fork.clone().unwrap_or_else(|| RepoPath {
373 namespace: repo.namespace.clone(),
374 name: repo.name.clone(),
375 });
376 let item = QueueStackItem {
377 number: pull.number,
378 title: pull.title.clone(),
379 source,
380 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
381 commit,
382 };
383 items.push((row, pull, item, checks));
384 }
385
386 // What the default branch has promised so far: every check an issue
387 // passed when it landed, newest first.
388 #[derive(Deserialize)]
389 struct ChecksRow {
390 checks: String,
391 }
392 let mut contract: Vec<String> = Vec::new();
393 for row in self
394 .db
395 .prepare(
396 "SELECT checks FROM issues
397 WHERE repo_id = ? AND state = 'closed' AND reason = 'completed' AND checks != '[]'
398 ORDER BY closed_at DESC LIMIT ?",
399 )
400 .bind(&[repo.id.as_str().into(), CONTRACT_ISSUES.into()])?
401 .all()
402 .await?
403 .results::<ChecksRow>()?
404 {
405 for check in serde_json::from_str::<Vec<String>>(&row.checks).unwrap_or_default() {
406 if !contract.contains(&check) {
407 contract.push(check);
408 }
409 }
410 }
411
412 let now = rfc3339(now_ms());
413 let mut jobs = Vec::new();
414 for index in 0..items.len() {
415 let (row, pull, _, _) = &items[index];
416 let stack: Vec<QueueStackItem> = items[..=index].iter().map(|(_, _, item, _)| item.clone()).collect();
417 let mut checks: Vec<String> = Vec::new();
418 for (_, _, _, theirs) in &items[..=index] {
419 for check in theirs {
420 if !checks.contains(check) {
421 checks.push(check.clone());
422 }
423 }
424 }
425 let ahead: Vec<u32> = stack[..index].iter().map(|item| item.number).collect();
426 let token = new_token();
427 self.db
428 .prepare(
429 "UPDATE queue_entries
430 SET state = 'testing', token_hash = ?, head_commit = ?, base_commit = ?,
431 ahead = ?, combined_commit = NULL, results = NULL, error = NULL,
432 tested_at = ?
433 WHERE id = ? AND state = 'waiting'",
434 )
435 .bind(&[
436 hash(&token).into(),
437 stack[index].commit.as_str().into(),
438 base.as_str().into(),
439 serde_json::to_string(&ahead)?.into(),
440 now.as_str().into(),
441 row.id.as_str().into(),
442 ])?
443 .run()
444 .await?;
445 // The sandbox reads the change and pushes the tested state as
446 // the pull request's author: a real account, whose token carries
447 // its memberships. (Whoever queued it may be g1t itself.)
448 let actor = pull.author.clone();
449 jobs.push(QueueJob {
450 entry_id: row.id.clone(),
451 token,
452 repo: RepoPath {
453 namespace: repo.namespace.clone(),
454 name: repo.name.clone(),
455 },
456 default_branch: repo.default_branch.clone(),
457 base_commit: base.clone(),
458 branch: row.branch(),
459 contract_checks: contract.iter().filter(|check| !checks.contains(check)).cloned().collect(),
460 stack,
461 checks,
462 actor,
463 });
464 }
465 Ok(jobs)
466 }
467
468 /// A sandbox's result for one combined state.
469 pub(crate) async fn report_queue(&self, a: ReportQueueArgs) -> Result<Outcome<QueueState>> {
470 let row = self
471 .db
472 .prepare("SELECT * FROM queue_entries WHERE id = ? AND token_hash = ?")
473 .bind(&[a.entry_id.as_str().into(), hash(&a.token).into()])?
474 .first::<EntryRow>(None)
475 .await?;
476 let Some(row) = row else {
477 return Ok(Outcome::fail(FailureCode::NotFound, "No such queue entry."));
478 };
479 if row.state() != QueueState::Testing {
480 return Ok(Outcome::fail(
481 FailureCode::Conflict,
482 "This state is no longer being tested.",
483 ));
484 }
485 let passed = a.error.is_none() && a.results.iter().all(|result| result.passed);
Sidebar: the panels really slide486 // A state that passed its checks still runs the repository's
487 // `merge_group` workflows, as GitHub's merge queue does; it stays in
488 // testing until they finish (see `statuses`).
489 let workflows = match (&a.combined_commit, passed) {
490 (Some(commit), true) => self.start_merge_group(&row, commit).await.unwrap_or(0),
491 _ => 0,
492 };
493 let state = if !passed {
494 QueueState::Failed
495 } else if workflows > 0 {
496 QueueState::Testing
497 } else {
498 QueueState::Passed
499 };
Agents as a team: lifecycle, merge queue, billing and a new shell500 self.db
501 .prepare(
502 "UPDATE queue_entries
503 SET state = ?, combined_commit = ?, results = ?, error = ?, token_hash = NULL,
504 finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END
505 WHERE id = ?",
506 )
507 .bind(&[
508 state.as_str().into(),
509 a.combined_commit.as_deref().map_or(JsValue::NULL, JsValue::from),
510 serde_json::to_string(&a.results)?.into(),
511 a.error.as_deref().map_or(JsValue::NULL, JsValue::from),
512 state.as_str().into(),
513 rfc3339(now_ms()).into(),
514 row.id.as_str().into(),
515 ])?
516 .run()
517 .await?;
518 if !passed {
519 self.eject(&row, &a).await?;
520 }
521 self.settle(&row.repo_id).await?;
522 self.changed(&row.repo_id).await?;
523 Ok(Outcome::Ok(state))
524 }
525
Sidebar: the panels really slide526 /// Asks the actions service to run the repository's `merge_group`
527 /// workflows on a combined state. Returns how many runs started.
528 async fn start_merge_group(&self, row: &EntryRow, commit: &str) -> Result<u32> {
529 #[derive(serde::Deserialize)]
530 struct Started {
531 runs: u32,
532 }
533 let started: Outcome<Started> = g1t_kit::call(
534 &self.actions,
535 "merge_group",
536 &serde_json::json!({
537 "repoId": row.repo_id,
538 "entry": row.id,
539 "sha": commit,
540 "headRef": format!("refs/heads/{}", row.branch()),
541 "baseSha": row.base_commit,
542 "number": row.number,
543 "ahead": row.ahead(),
544 }),
545 )
546 .await?;
547 Ok(match started {
548 Outcome::Ok(started) => started.runs,
549 Outcome::Fail(_) => 0,
550 })
551 }
552
553 /// Workflows on a combined state finished: it passes and lands in turn,
554 /// or fails and leaves the queue as for failed checks.
555 pub(crate) async fn merge_group_finished(&self, repo_id: &str, commit: &str, failed: &[String]) -> Result<()> {
556 let row = self
557 .db
558 .prepare(
559 "SELECT * FROM queue_entries WHERE repo_id = ? AND combined_commit = ? AND state = 'testing' AND token_hash IS NULL",
560 )
561 .bind(&[repo_id.into(), commit.into()])?
562 .first::<EntryRow>(None)
563 .await?;
564 let Some(row) = row else { return Ok(()) };
565 let passed = failed.is_empty();
566 self.db
567 .prepare("UPDATE queue_entries SET state = ?, finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END WHERE id = ?")
568 .bind(&[
569 (if passed { "passed" } else { "failed" }).into(),
570 (if passed { "passed" } else { "failed" }).into(),
571 rfc3339(now_ms()).into(),
572 row.id.as_str().into(),
573 ])?
574 .run()
575 .await?;
576 if !passed {
577 let report = ReportQueueArgs {
578 entry_id: row.id.clone(),
579 token: String::new(),
580 combined_commit: Some(commit.to_owned()),
581 results: Vec::new(),
582 error: Some(format!(
583 "the workflow {} failed on it",
584 crate::statuses::list(failed)
585 )),
586 conflict_with: None,
587 };
588 self.eject(&row, &report).await?;
589 }
590 self.settle(repo_id).await?;
591 self.changed(repo_id).await?;
592 Ok(())
593 }
594
Agents as a team: lifecycle, merge queue, billing and a new shell595 /// An entry whose combined state failed: the entries tested on top of
596 /// it are tested again without it, and the failure is recorded as a
597 /// failed check run of its pull request, so a g1t agent is sent back.
598 async fn eject(&self, row: &EntryRow, report: &ReportQueueArgs) -> Result<()> {
599 let active = self.entries(&row.repo_id, true).await?;
600 let behind: Vec<&EntryRow> = active
601 .iter()
602 .filter(|other| other.state() != QueueState::Waiting && other.ahead().contains(&row.number))
603 .collect();
604 self.retest(&behind).await?;
605
606 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
607 return Ok(());
608 };
609 let ahead = row.ahead();
610 let state = if ahead.is_empty() {
611 "the default branch as it is now".to_owned()
612 } else {
613 format!(
614 "the default branch with {} merged in first",
615 ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ")
616 )
617 };
618 let why = match (&report.error, report.conflict_with) {
619 (_, Some(other)) if other != row.number => format!(
620 "Its change conflicts with #{other}, which is ahead of it in the merge queue. Bring it up to date with the default branch once #{other} lands, and merge it again."
621 ),
Sidebar: the panels really slide622 (Some(error), _) if error.starts_with("the workflow ") => {
623 format!("{} when it was combined with {state}.", error.replacen("the workflow", "The workflow", 1).trim_end_matches(" on it"))
624 }
Agents as a team: lifecycle, merge queue, billing and a new shell625 (Some(error), _) => format!("Its combined state could not be built or checked: {error}"),
626 (None, _) => format!(
627 "The acceptance checks failed when it was combined with {state}, though it may pass on its own."
628 ),
629 };
630 let now = now_ms();
631 let run_id = new_id("chk", now);
632 let mut results = report.results.clone();
633 for result in &mut results {
634 result.command = format!("{} (merge queue, on {state})", result.command);
635 }
636 self.db
637 .batch(vec![
638 self.db
639 .prepare(
640 "INSERT INTO check_runs
641 (id, pull_id, head_commit, status, results, error, token_hash,
642 created_at, finished_at)
643 VALUES (?, ?, ?, 'failed', ?, ?, ?, ?, ?)",
644 )
645 .bind(&[
646 run_id.as_str().into(),
647 pull.id.as_str().into(),
648 row.head_commit.as_deref().unwrap_or_default().into(),
649 serde_json::to_string(&results)?.into(),
650 why.as_str().into(),
651 hash(&new_token()).into(),
652 rfc3339(now).into(),
653 rfc3339(now).into(),
654 ])?,
655 self.db
656 .prepare("UPDATE pulls SET check_status = 'failed', check_run_id = ? WHERE id = ?")
657 .bind(&[run_id.as_str().into(), pull.id.as_str().into()])?,
658 ])
659 .await?;
660 self.note(
661 &row.repo_id,
662 pull.number,
663 ("g1t", "g1t"),
664 &format!("was taken out of the merge queue. {why}"),
665 )
666 .await?;
667 self.publish_as(
668 "checks.completed",
669 &row.repo_id,
670 None,
671 ChecksEvent {
672 pull_id: pull.id.clone(),
673 repo_id: row.repo_id.clone(),
674 number: pull.number,
675 status: "failed",
676 commit: row.head_commit.clone().unwrap_or_default(),
677 },
678 )
679 .await?;
680 Ok(())
681 }
682
683 /// Lands every entry at the front of the queue whose tested state
684 /// passed, in order.
685 async fn settle(&self, repo_id: &str) -> Result<()> {
686 let active = self.entries(repo_id, true).await?;
687 let Some(viewer) = active.first().and_then(EntryRow::actor) else {
688 return Ok(());
689 };
690 let repo: Outcome<Repo> = g1t_kit::call(
691 &self.repos,
692 "get_by_id",
693 &GetByIdArgs {
694 id: repo_id.to_owned(),
695 viewer: Some(viewer),
696 },
697 )
698 .await?;
699 let Outcome::Ok(repo) = repo else {
700 return Ok(());
701 };
702 for row in active {
703 if row.state() != QueueState::Passed {
704 // The front has not passed yet: nothing behind it may land.
705 return Ok(());
706 }
707 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
708 continue;
709 };
710 let Some(actor) = row.actor() else {
711 continue;
712 };
713 // Pushed to since it was tested: test it again as it is now.
714 if pull.status != PullStatus::Open {
715 self.leave(repo_id, &pull, QueueState::Removed, Some("It was closed."))
716 .await?;
717 continue;
718 }
719 let head: Option<String> = g1t_kit::call(
720 &self.repos,
721 "head",
722 &HeadArgs {
723 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo_id.to_owned()),
724 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
725 },
726 )
727 .await?;
728 if head.is_some() && head != row.head_commit {
729 let active = self.entries(repo_id, true).await?;
730 let again: Vec<&EntryRow> = active
731 .iter()
732 .filter(|other| other.id == row.id || other.ahead().contains(&row.number))
733 .collect();
734 self.retest(&again).await?;
735 return Ok(());
736 }
737 let landed: Outcome<Landed> = g1t_kit::call(
738 &self.repos,
739 "land",
740 &LandArgs {
741 source_id: repo_id.to_owned(),
742 branch: Some(row.branch()),
743 actor: actor.clone(),
744 },
745 )
746 .await?;
747 let landed = match landed {
748 Outcome::Ok(landed) => landed,
749 Outcome::Fail(_) => {
750 // The default branch moved outside the queue: every
751 // tested state is built on something that is gone.
752 let active = self.entries(repo_id, true).await?;
753 let all: Vec<&EntryRow> = active.iter().collect();
754 self.retest(&all).await?;
755 return Ok(());
756 }
757 };
758 self.db
759 .prepare(
760 "UPDATE queue_entries SET state = 'landed', finished_at = ? WHERE id = ?",
761 )
762 .bind(&[rfc3339(now_ms()).into(), row.id.as_str().into()])?
763 .run()
764 .await?;
765 self.record_merge(&repo, pull, &actor, row.keep_issue_open != 0, landed)
766 .await?;
767 }
768 Ok(())
769 }
770}