Skip to content

g1t/services/work/src/queue.rs

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

This file's history is long; its oldest lines are credited to the oldest commit read.