flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/queue.rs

800 lines31,932 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 author: a real account, whose token carries
434 // its memberships. (Whoever queued it may be g1t itself.)
435 let actor = pull.author.clone();
436 jobs.push(QueueJob {
437 entry_id: row.id.clone(),
438 token,
439 repo: RepoPath {
440 namespace: repo.namespace.clone(),
441 name: repo.name.clone(),
442 },
443 default_branch: repo.default_branch.clone(),
444 base_commit: base.clone(),
445 branch: row.branch(),
446 // The state is checked by its merge_group workflows.
447 contract_checks: Vec::new(),
448 stack,
449 checks: Vec::new(),
450 actor,
451 });
452 }
453 Ok(jobs)
454 }
455
456 /// A sandbox's result for one combined state.
457 pub(crate) async fn report_queue(&self, a: ReportQueueArgs) -> Result<Outcome<QueueState>> {
458 let row = self
459 .db
460 .prepare("SELECT * FROM queue_entries WHERE id = ? AND token_hash = ?")
461 .bind(&[a.entry_id.as_str().into(), hash(&a.token).into()])?
462 .first::<EntryRow>(None)
463 .await?;
464 let Some(row) = row else {
465 return Ok(Outcome::fail(FailureCode::NotFound, "No such queue entry."));
466 };
467 if row.state() != QueueState::Testing {
468 return Ok(Outcome::fail(
469 FailureCode::Conflict,
470 "This state is no longer being tested.",
471 ));
472 }
473 let built = a.error.is_none() && a.results.iter().all(|result| result.passed);
474 // A state that was built runs the repository's `merge_group`
475 // workflows; it stays in testing until they finish (see `statuses`),
476 // and the branch's required checks must pass on it.
477 let workflows = match (&a.combined_commit, built) {
478 (Some(commit), true) => self.start_merge_group(&row, commit).await.unwrap_or(0),
479 _ => 0,
480 };
481 let required = self.settings(&row.repo_id).await?.required_checks;
482 // Nothing runs on it, so the required checks never would report.
483 let unchecked = (built && workflows == 0 && !required.is_empty()).then(|| {
484 format!(
485 "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",
486 if required.len() == 1 { "check" } else { "checks" },
487 crate::statuses::list(&required)
488 )
489 });
490 let passed = built && unchecked.is_none();
491 let error = a.error.clone().or(unchecked);
492 let state = if !passed {
493 QueueState::Failed
494 } else if workflows > 0 {
495 QueueState::Testing
496 } else {
497 QueueState::Passed
498 };
499 self.db
500 .prepare(
501 "UPDATE queue_entries
502 SET state = ?, combined_commit = ?, results = ?, error = ?, token_hash = NULL,
503 finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END
504 WHERE id = ?",
505 )
506 .bind(&[
507 state.as_str().into(),
508 a.combined_commit.as_deref().map_or(JsValue::NULL, JsValue::from),
509 serde_json::to_string(&a.results)?.into(),
510 error.as_deref().map_or(JsValue::NULL, JsValue::from),
511 state.as_str().into(),
512 rfc3339(now_ms()).into(),
513 row.id.as_str().into(),
514 ])?
515 .run()
516 .await?;
517 if !passed {
518 let report = ReportQueueArgs { error, ..a };
519 self.eject(&row, &report).await?;
520 }
521 self.settle(&row.repo_id).await?;
522 self.changed(&row.repo_id).await?;
523 Ok(Outcome::Ok(state))
524 }
525
526 /// 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 /// when they all passed and so did every required check, or fails and
555 /// leaves the queue.
556 pub(crate) async fn merge_group_finished(
557 &self,
558 repo_id: &str,
559 commit: &str,
560 facts: &crate::statuses::WorkflowFacts,
561 ) -> Result<()> {
562 let row = self
563 .db
564 .prepare(
565 "SELECT * FROM queue_entries WHERE repo_id = ? AND combined_commit = ? AND state = 'testing' AND token_hash IS NULL",
566 )
567 .bind(&[repo_id.into(), commit.into()])?
568 .first::<EntryRow>(None)
569 .await?;
570 let Some(row) = row else { return Ok(()) };
571 let missing = facts.expected();
572 let passed = facts.failed.is_empty() && facts.required_failed().is_empty() && missing.is_empty();
573 self.db
574 .prepare("UPDATE queue_entries SET state = ?, finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END WHERE id = ?")
575 .bind(&[
576 (if passed { "passed" } else { "failed" }).into(),
577 (if passed { "passed" } else { "failed" }).into(),
578 rfc3339(now_ms()).into(),
579 row.id.as_str().into(),
580 ])?
581 .run()
582 .await?;
583 if !passed {
584 let report = ReportQueueArgs {
585 entry_id: row.id.clone(),
586 token: String::new(),
587 combined_commit: Some(commit.to_owned()),
588 results: Vec::new(),
589 error: Some(if facts.failed.is_empty() {
590 format!(
591 "the required {} {} did not report on it. Add merge_group to the on: of the workflows the branch requires",
592 if missing.len() == 1 { "check" } else { "checks" },
593 crate::statuses::list(&missing)
594 )
595 } else {
596 format!("the workflow {} failed on it", crate::statuses::list(&facts.failed))
597 }),
598 conflict_with: None,
599 conflicts: Vec::new(),
600 };
601 self.eject(&row, &report).await?;
602 }
603 self.settle(repo_id).await?;
604 self.changed(repo_id).await?;
605 Ok(())
606 }
607
608 /// An entry whose combined state failed: the entries tested on top of
609 /// it are tested again without it, and the failure is recorded as a
610 /// failed check run of its pull request, so a g1t agent is sent back.
611 async fn eject(&self, row: &EntryRow, report: &ReportQueueArgs) -> Result<()> {
612 self.drop_branch(row).await;
613 let active = self.entries(&row.repo_id, true).await?;
614 let behind: Vec<&EntryRow> = active
615 .iter()
616 .filter(|other| other.state() != QueueState::Waiting && other.ahead().contains(&row.number))
617 .collect();
618 self.retest(&behind).await?;
619
620 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
621 return Ok(());
622 };
623 let ahead = row.ahead();
624 let state = if ahead.is_empty() {
625 "the default branch as it is now".to_owned()
626 } else {
627 format!(
628 "the default branch with {} merged in first",
629 ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ")
630 )
631 };
632 // Named in backticks, so the conversation can link each to the diff.
633 let files = crate::mergeability::tidy(report.conflicts.clone())
634 .iter()
635 .map(|path| format!("`{path}`"))
636 .collect::<Vec<_>>()
637 .join(", ");
638 let why = match (&report.error, report.conflict_with) {
639 (_, Some(other)) if other != row.number && !files.is_empty() => format!(
640 "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."
641 ),
642 (_, Some(other)) if other != row.number => format!(
643 "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."
644 ),
645 (Some(_), None) if !files.is_empty() => format!(
646 "Its change conflicts with the default branch in {files}. Bring it up to date with the default branch, and merge it again."
647 ),
648 (Some(error), _) if error.starts_with("the workflow ") => {
649 format!("{} when it was combined with {state}.", error.replacen("the workflow", "The workflow", 1).trim_end_matches(" on it"))
650 }
651 (Some(error), _) if error.starts_with("the required ") => {
652 format!("Combined with {state}, {error}.")
653 }
654 (Some(error), _) => format!("Its combined state could not be built or checked: {error}"),
655 (None, _) => format!(
656 "It failed when it was combined with {state}, though it may pass on its own."
657 ),
658 };
659 let now = now_ms();
660 let run_id = new_id("chk", now);
661 let mut results = report.results.clone();
662 for result in &mut results {
663 result.command = format!("{} (merge queue, on {state})", result.command);
664 }
665 self.db
666 .batch(vec![
667 self.db
668 .prepare(
669 "INSERT INTO check_runs
670 (id, pull_id, head_commit, status, results, error, token_hash,
671 created_at, finished_at)
672 VALUES (?, ?, ?, 'failed', ?, ?, ?, ?, ?)",
673 )
674 .bind(&[
675 run_id.as_str().into(),
676 pull.id.as_str().into(),
677 row.head_commit.as_deref().unwrap_or_default().into(),
678 serde_json::to_string(&results)?.into(),
679 why.as_str().into(),
680 hash(&new_token()).into(),
681 rfc3339(now).into(),
682 rfc3339(now).into(),
683 ])?,
684 self.db
685 .prepare("UPDATE pulls SET check_status = 'failed', check_run_id = ? WHERE id = ?")
686 .bind(&[run_id.as_str().into(), pull.id.as_str().into()])?,
687 ])
688 .await?;
689 self.note(
690 &row.repo_id,
691 pull.number,
692 ("g1t", "g1t"),
693 &format!("was taken out of the merge queue. {why}"),
694 )
695 .await?;
696 self.publish_as(
697 "checks.completed",
698 &row.repo_id,
699 None,
700 ChecksEvent {
701 pull_id: pull.id.clone(),
702 repo_id: row.repo_id.clone(),
703 number: pull.number,
704 status: "failed",
705 commit: row.head_commit.clone().unwrap_or_default(),
706 },
707 )
708 .await?;
709 Ok(())
710 }
711
712 /// Lands every entry at the front of the queue whose tested state
713 /// passed, in order.
714 async fn settle(&self, repo_id: &str) -> Result<()> {
715 let active = self.entries(repo_id, true).await?;
716 let Some(viewer) = active.first().and_then(EntryRow::actor) else {
717 return Ok(());
718 };
719 let repo: Outcome<Repo> = g1t_kit::call(
720 &self.repos,
721 "get_by_id",
722 &GetByIdArgs {
723 id: repo_id.to_owned(),
724 viewer: Some(viewer),
725 },
726 )
727 .await?;
728 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
729 return Ok(());
730 };
731 for row in active {
732 if row.state() != QueueState::Passed {
733 // The front has not passed yet: nothing behind it may land.
734 return Ok(());
735 }
736 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
737 continue;
738 };
739 let Some(actor) = row.actor() else {
740 continue;
741 };
742 // Pushed to since it was tested: test it again as it is now.
743 if pull.status != PullStatus::Open {
744 self.leave(repo_id, &pull, QueueState::Removed, Some("It was closed."))
745 .await?;
746 continue;
747 }
748 let head: Option<String> = g1t_kit::call(
749 &self.repos,
750 "head",
751 &HeadArgs {
752 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo_id.to_owned()),
753 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
754 },
755 )
756 .await?;
757 if head.is_some() && head != row.head_commit {
758 let active = self.entries(repo_id, true).await?;
759 let again: Vec<&EntryRow> = active
760 .iter()
761 .filter(|other| other.id == row.id || other.ahead().contains(&row.number))
762 .collect();
763 self.retest(&again).await?;
764 return Ok(());
765 }
766 let landed: Outcome<Landed> = g1t_kit::call(
767 &self.repos,
768 "land",
769 &LandArgs {
770 source_id: repo_id.to_owned(),
771 branch: Some(row.branch()),
772 actor: actor.clone(),
773 },
774 )
775 .await?;
776 let landed = match landed {
777 Outcome::Ok(landed) => landed,
778 Outcome::Fail(_) => {
779 // The default branch moved outside the queue: every
780 // tested state is built on something that is gone.
781 let active = self.entries(repo_id, true).await?;
782 let all: Vec<&EntryRow> = active.iter().collect();
783 self.retest(&all).await?;
784 return Ok(());
785 }
786 };
787 self.db
788 .prepare(
789 "UPDATE queue_entries SET state = 'landed', finished_at = ? WHERE id = ?",
790 )
791 .bind(&[rfc3339(now_ms()).into(), row.id.as_str().into()])?
792 .run()
793 .await?;
794 self.drop_branch(&row).await;
795 self.record_merge(&repo, pull, &actor, row.keep_issue_open != 0, landed)
796 .await?;
797 }
798 Ok(())
799 }
800}