g1t/services/work/src/queue.rs

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