pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/queue.rs

793 lines30,809 bytesCodeBlame
1//! 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::{DeleteBranchArgs, 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)]
39pub(crate) struct EntryRow {
40 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 self.drop_branch(entry).await;
264 let behind: Vec<&EntryRow> = active
265 .iter()
266 .filter(|row| row.state() != QueueState::Waiting && row.ahead().contains(&pull.number))
267 .collect();
268 self.retest(&behind).await?;
269 Ok(true)
270 }
271
272 /// Removes an entry's tested state from the repository once it has
273 /// left the queue, as GitHub does with its queue's branches. A branch
274 /// left behind is untidy, not wrong, so a failure only logs.
275 async fn drop_branch(&self, row: &EntryRow) {
276 let deleted: Result<Outcome<bool>> = g1t_kit::call(
277 &self.repos,
278 "delete_branch",
279 &DeleteBranchArgs {
280 repo_id: row.repo_id.clone(),
281 branch: row.branch(),
282 },
283 )
284 .await;
285 match deleted {
286 Ok(Outcome::Ok(_)) => {}
287 Ok(Outcome::Fail(failure)) => worker::console_warn!("{}: {}", row.branch(), failure.message),
288 Err(error) => worker::console_warn!("{}: {error}", row.branch()),
289 }
290 }
291
292 /// Sends entries back to waiting, to be tested again.
293 async fn retest(&self, rows: &[&EntryRow]) -> Result<()> {
294 for row in rows {
295 self.db
296 .prepare(
297 "UPDATE queue_entries
298 SET state = 'waiting', token_hash = NULL, combined_commit = NULL,
299 results = NULL, error = NULL, ahead = NULL
300 WHERE id = ? AND state IN ('testing', 'passed')",
301 )
302 .bind(&[row.id.as_str().into()])?
303 .run()
304 .await?;
305 }
306 Ok(())
307 }
308
309 /// The next batch to test, if nothing is being tested now: one job per
310 /// entry, each building the default branch with that entry and every
311 /// entry ahead of it.
312 pub(crate) async fn queue_build(&self, a: QueueBuildArgs) -> Result<Vec<QueueJob>> {
313 let active = self.entries(&a.repo_id, true).await?;
314 // A batch that has taken too long is tested again.
315 let stale = minutes_ago(TESTING_MINUTES);
316 let stuck: Vec<&EntryRow> = active
317 .iter()
318 .filter(|row| {
319 row.state() == QueueState::Testing
320 && row.tested_at.as_deref().is_none_or(|at| at < stale.as_str())
321 })
322 .collect();
323 if !stuck.is_empty() {
324 self.retest(&stuck).await?;
325 return Box::pin(self.queue_build(a)).await;
326 }
327 if active.iter().any(|row| row.state() != QueueState::Waiting) {
328 return Ok(Vec::new());
329 }
330 let batch: Vec<EntryRow> = active.into_iter().take(BATCH).collect();
331 let Some(first) = batch.first() else {
332 return Ok(Vec::new());
333 };
334 let Some(actor) = first.actor() else {
335 return Ok(Vec::new());
336 };
337 let repo: Outcome<Repo> = g1t_kit::call(
338 &self.repos,
339 "get_by_id",
340 &GetByIdArgs {
341 id: a.repo_id.clone(),
342 viewer: Some(actor),
343 },
344 )
345 .await?;
346 let Outcome::Ok(repo) = repo else {
347 return Ok(Vec::new());
348 };
349 let base: Option<String> = g1t_kit::call(
350 &self.repos,
351 "head",
352 &HeadArgs {
353 repo_id: repo.id.clone(),
354 branch: repo.default_branch.clone(),
355 },
356 )
357 .await?;
358 let Some(base) = base else {
359 return Ok(Vec::new());
360 };
361
362 // Each entry's change, as it is now, and the checks it brings.
363 let mut items: Vec<(EntryRow, Pull, QueueStackItem, Vec<String>)> = Vec::new();
364 for row in batch {
365 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
366 continue;
367 };
368 if pull.status != PullStatus::Open {
369 self.leave(&repo.id, &pull, QueueState::Removed, Some("It was closed."))
370 .await?;
371 continue;
372 }
373 let head: Option<String> = g1t_kit::call(
374 &self.repos,
375 "head",
376 &HeadArgs {
377 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
378 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
379 },
380 )
381 .await?;
382 let Some(commit) = head else {
383 continue;
384 };
385 let checks = match pull.issue {
386 Some(number) => self
387 .issue(&repo.id, number)
388 .await?
389 .map(|issue| issue.checks)
390 .unwrap_or_default(),
391 None => Vec::new(),
392 };
393 let source = pull.fork.clone().unwrap_or_else(|| RepoPath {
394 namespace: repo.namespace.clone(),
395 name: repo.name.clone(),
396 });
397 let item = QueueStackItem {
398 number: pull.number,
399 title: pull.title.clone(),
400 source,
401 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
402 commit,
403 };
404 items.push((row, pull, item, checks));
405 }
406
407 // What the default branch has promised so far: every check an issue
408 // passed when it landed, newest first.
409 #[derive(Deserialize)]
410 struct ChecksRow {
411 checks: String,
412 }
413 let mut contract: Vec<String> = Vec::new();
414 for row in self
415 .db
416 .prepare(
417 "SELECT checks FROM issues
418 WHERE repo_id = ? AND state = 'closed' AND reason = 'completed' AND checks != '[]'
419 ORDER BY closed_at DESC LIMIT ?",
420 )
421 .bind(&[repo.id.as_str().into(), CONTRACT_ISSUES.into()])?
422 .all()
423 .await?
424 .results::<ChecksRow>()?
425 {
426 for check in serde_json::from_str::<Vec<String>>(&row.checks).unwrap_or_default() {
427 if !contract.contains(&check) {
428 contract.push(check);
429 }
430 }
431 }
432
433 let now = rfc3339(now_ms());
434 let mut jobs = Vec::new();
435 for index in 0..items.len() {
436 let (row, pull, _, _) = &items[index];
437 let stack: Vec<QueueStackItem> = items[..=index].iter().map(|(_, _, item, _)| item.clone()).collect();
438 let mut checks: Vec<String> = Vec::new();
439 for (_, _, _, theirs) in &items[..=index] {
440 for check in theirs {
441 if !checks.contains(check) {
442 checks.push(check.clone());
443 }
444 }
445 }
446 let ahead: Vec<u32> = stack[..index].iter().map(|item| item.number).collect();
447 let token = new_token();
448 self.db
449 .prepare(
450 "UPDATE queue_entries
451 SET state = 'testing', token_hash = ?, head_commit = ?, base_commit = ?,
452 ahead = ?, combined_commit = NULL, results = NULL, error = NULL,
453 tested_at = ?
454 WHERE id = ? AND state = 'waiting'",
455 )
456 .bind(&[
457 hash(&token).into(),
458 stack[index].commit.as_str().into(),
459 base.as_str().into(),
460 serde_json::to_string(&ahead)?.into(),
461 now.as_str().into(),
462 row.id.as_str().into(),
463 ])?
464 .run()
465 .await?;
466 // The sandbox reads the change and pushes the tested state as
467 // the pull request's author: a real account, whose token carries
468 // its memberships. (Whoever queued it may be g1t itself.)
469 let actor = pull.author.clone();
470 jobs.push(QueueJob {
471 entry_id: row.id.clone(),
472 token,
473 repo: RepoPath {
474 namespace: repo.namespace.clone(),
475 name: repo.name.clone(),
476 },
477 default_branch: repo.default_branch.clone(),
478 base_commit: base.clone(),
479 branch: row.branch(),
480 contract_checks: contract.iter().filter(|check| !checks.contains(check)).cloned().collect(),
481 stack,
482 checks,
483 actor,
484 });
485 }
486 Ok(jobs)
487 }
488
489 /// A sandbox's result for one combined state.
490 pub(crate) async fn report_queue(&self, a: ReportQueueArgs) -> Result<Outcome<QueueState>> {
491 let row = self
492 .db
493 .prepare("SELECT * FROM queue_entries WHERE id = ? AND token_hash = ?")
494 .bind(&[a.entry_id.as_str().into(), hash(&a.token).into()])?
495 .first::<EntryRow>(None)
496 .await?;
497 let Some(row) = row else {
498 return Ok(Outcome::fail(FailureCode::NotFound, "No such queue entry."));
499 };
500 if row.state() != QueueState::Testing {
501 return Ok(Outcome::fail(
502 FailureCode::Conflict,
503 "This state is no longer being tested.",
504 ));
505 }
506 let passed = a.error.is_none() && a.results.iter().all(|result| result.passed);
507 // A state that passed its checks still runs the repository's
508 // `merge_group` workflows, as GitHub's merge queue does; it stays in
509 // testing until they finish (see `statuses`).
510 let workflows = match (&a.combined_commit, passed) {
511 (Some(commit), true) => self.start_merge_group(&row, commit).await.unwrap_or(0),
512 _ => 0,
513 };
514 let state = if !passed {
515 QueueState::Failed
516 } else if workflows > 0 {
517 QueueState::Testing
518 } else {
519 QueueState::Passed
520 };
521 self.db
522 .prepare(
523 "UPDATE queue_entries
524 SET state = ?, combined_commit = ?, results = ?, error = ?, token_hash = NULL,
525 finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END
526 WHERE id = ?",
527 )
528 .bind(&[
529 state.as_str().into(),
530 a.combined_commit.as_deref().map_or(JsValue::NULL, JsValue::from),
531 serde_json::to_string(&a.results)?.into(),
532 a.error.as_deref().map_or(JsValue::NULL, JsValue::from),
533 state.as_str().into(),
534 rfc3339(now_ms()).into(),
535 row.id.as_str().into(),
536 ])?
537 .run()
538 .await?;
539 if !passed {
540 self.eject(&row, &a).await?;
541 }
542 self.settle(&row.repo_id).await?;
543 self.changed(&row.repo_id).await?;
544 Ok(Outcome::Ok(state))
545 }
546
547 /// Asks the actions service to run the repository's `merge_group`
548 /// workflows on a combined state. Returns how many runs started.
549 async fn start_merge_group(&self, row: &EntryRow, commit: &str) -> Result<u32> {
550 #[derive(serde::Deserialize)]
551 struct Started {
552 runs: u32,
553 }
554 let started: Outcome<Started> = g1t_kit::call(
555 &self.actions,
556 "merge_group",
557 &serde_json::json!({
558 "repoId": row.repo_id,
559 "entry": row.id,
560 "sha": commit,
561 "headRef": format!("refs/heads/{}", row.branch()),
562 "baseSha": row.base_commit,
563 "number": row.number,
564 "ahead": row.ahead(),
565 }),
566 )
567 .await?;
568 Ok(match started {
569 Outcome::Ok(started) => started.runs,
570 Outcome::Fail(_) => 0,
571 })
572 }
573
574 /// Workflows on a combined state finished: it passes and lands in turn,
575 /// or fails and leaves the queue as for failed checks.
576 pub(crate) async fn merge_group_finished(&self, repo_id: &str, commit: &str, failed: &[String]) -> Result<()> {
577 let row = self
578 .db
579 .prepare(
580 "SELECT * FROM queue_entries WHERE repo_id = ? AND combined_commit = ? AND state = 'testing' AND token_hash IS NULL",
581 )
582 .bind(&[repo_id.into(), commit.into()])?
583 .first::<EntryRow>(None)
584 .await?;
585 let Some(row) = row else { return Ok(()) };
586 let passed = failed.is_empty();
587 self.db
588 .prepare("UPDATE queue_entries SET state = ?, finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END WHERE id = ?")
589 .bind(&[
590 (if passed { "passed" } else { "failed" }).into(),
591 (if passed { "passed" } else { "failed" }).into(),
592 rfc3339(now_ms()).into(),
593 row.id.as_str().into(),
594 ])?
595 .run()
596 .await?;
597 if !passed {
598 let report = ReportQueueArgs {
599 entry_id: row.id.clone(),
600 token: String::new(),
601 combined_commit: Some(commit.to_owned()),
602 results: Vec::new(),
603 error: Some(format!(
604 "the workflow {} failed on it",
605 crate::statuses::list(failed)
606 )),
607 conflict_with: None,
608 };
609 self.eject(&row, &report).await?;
610 }
611 self.settle(repo_id).await?;
612 self.changed(repo_id).await?;
613 Ok(())
614 }
615
616 /// An entry whose combined state failed: the entries tested on top of
617 /// it are tested again without it, and the failure is recorded as a
618 /// failed check run of its pull request, so a g1t agent is sent back.
619 async fn eject(&self, row: &EntryRow, report: &ReportQueueArgs) -> Result<()> {
620 self.drop_branch(row).await;
621 let active = self.entries(&row.repo_id, true).await?;
622 let behind: Vec<&EntryRow> = active
623 .iter()
624 .filter(|other| other.state() != QueueState::Waiting && other.ahead().contains(&row.number))
625 .collect();
626 self.retest(&behind).await?;
627
628 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
629 return Ok(());
630 };
631 let ahead = row.ahead();
632 let state = if ahead.is_empty() {
633 "the default branch as it is now".to_owned()
634 } else {
635 format!(
636 "the default branch with {} merged in first",
637 ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ")
638 )
639 };
640 let why = match (&report.error, report.conflict_with) {
641 (_, Some(other)) if other != row.number => format!(
642 "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."
643 ),
644 (Some(error), _) if error.starts_with("the workflow ") => {
645 format!("{} when it was combined with {state}.", error.replacen("the workflow", "The workflow", 1).trim_end_matches(" on it"))
646 }
647 (Some(error), _) => format!("Its combined state could not be built or checked: {error}"),
648 (None, _) => format!(
649 "The acceptance checks failed when it was combined with {state}, though it may pass on its own."
650 ),
651 };
652 let now = now_ms();
653 let run_id = new_id("chk", now);
654 let mut results = report.results.clone();
655 for result in &mut results {
656 result.command = format!("{} (merge queue, on {state})", result.command);
657 }
658 self.db
659 .batch(vec![
660 self.db
661 .prepare(
662 "INSERT INTO check_runs
663 (id, pull_id, head_commit, status, results, error, token_hash,
664 created_at, finished_at)
665 VALUES (?, ?, ?, 'failed', ?, ?, ?, ?, ?)",
666 )
667 .bind(&[
668 run_id.as_str().into(),
669 pull.id.as_str().into(),
670 row.head_commit.as_deref().unwrap_or_default().into(),
671 serde_json::to_string(&results)?.into(),
672 why.as_str().into(),
673 hash(&new_token()).into(),
674 rfc3339(now).into(),
675 rfc3339(now).into(),
676 ])?,
677 self.db
678 .prepare("UPDATE pulls SET check_status = 'failed', check_run_id = ? WHERE id = ?")
679 .bind(&[run_id.as_str().into(), pull.id.as_str().into()])?,
680 ])
681 .await?;
682 self.note(
683 &row.repo_id,
684 pull.number,
685 ("g1t", "g1t"),
686 &format!("was taken out of the merge queue. {why}"),
687 )
688 .await?;
689 self.publish_as(
690 "checks.completed",
691 &row.repo_id,
692 None,
693 ChecksEvent {
694 pull_id: pull.id.clone(),
695 repo_id: row.repo_id.clone(),
696 number: pull.number,
697 status: "failed",
698 commit: row.head_commit.clone().unwrap_or_default(),
699 },
700 )
701 .await?;
702 Ok(())
703 }
704
705 /// Lands every entry at the front of the queue whose tested state
706 /// passed, in order.
707 async fn settle(&self, repo_id: &str) -> Result<()> {
708 let active = self.entries(repo_id, true).await?;
709 let Some(viewer) = active.first().and_then(EntryRow::actor) else {
710 return Ok(());
711 };
712 let repo: Outcome<Repo> = g1t_kit::call(
713 &self.repos,
714 "get_by_id",
715 &GetByIdArgs {
716 id: repo_id.to_owned(),
717 viewer: Some(viewer),
718 },
719 )
720 .await?;
721 let Outcome::Ok(repo) = repo else {
722 return Ok(());
723 };
724 for row in active {
725 if row.state() != QueueState::Passed {
726 // The front has not passed yet: nothing behind it may land.
727 return Ok(());
728 }
729 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
730 continue;
731 };
732 let Some(actor) = row.actor() else {
733 continue;
734 };
735 // Pushed to since it was tested: test it again as it is now.
736 if pull.status != PullStatus::Open {
737 self.leave(repo_id, &pull, QueueState::Removed, Some("It was closed."))
738 .await?;
739 continue;
740 }
741 let head: Option<String> = g1t_kit::call(
742 &self.repos,
743 "head",
744 &HeadArgs {
745 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo_id.to_owned()),
746 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
747 },
748 )
749 .await?;
750 if head.is_some() && head != row.head_commit {
751 let active = self.entries(repo_id, true).await?;
752 let again: Vec<&EntryRow> = active
753 .iter()
754 .filter(|other| other.id == row.id || other.ahead().contains(&row.number))
755 .collect();
756 self.retest(&again).await?;
757 return Ok(());
758 }
759 let landed: Outcome<Landed> = g1t_kit::call(
760 &self.repos,
761 "land",
762 &LandArgs {
763 source_id: repo_id.to_owned(),
764 branch: Some(row.branch()),
765 actor: actor.clone(),
766 },
767 )
768 .await?;
769 let landed = match landed {
770 Outcome::Ok(landed) => landed,
771 Outcome::Fail(_) => {
772 // The default branch moved outside the queue: every
773 // tested state is built on something that is gone.
774 let active = self.entries(repo_id, true).await?;
775 let all: Vec<&EntryRow> = active.iter().collect();
776 self.retest(&all).await?;
777 return Ok(());
778 }
779 };
780 self.db
781 .prepare(
782 "UPDATE queue_entries SET state = 'landed', finished_at = ? WHERE id = ?",
783 )
784 .bind(&[rfc3339(now_ms()).into(), row.id.as_str().into()])?
785 .run()
786 .await?;
787 self.drop_branch(&row).await;
788 self.record_merge(&repo, pull, &actor, row.keep_issue_open != 0, landed)
789 .await?;
790 }
791 Ok(())
792 }
793}