Skip to content
809 linesCodeBlameRaw

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.

Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request1//! 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
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request6//! 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.
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request14
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};
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request18use 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 {
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request36 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 };
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request120 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;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request178 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.
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request215 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));
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request231 }
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;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request269 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(),
Merge update PRs close themselves: g1t closes its security and version updates once they are no longer needed, and deletes their branches287 head: None,
Merge queue: tested states are deleted once their entry leaves288 },
289 )
290 .await;
291 match deleted {
292 Ok(Outcome::Ok(_)) => {}
293 Ok(Outcome::Fail(failure)) => worker::console_warn!("{}: {}", row.branch(), failure.message),
294 Err(error) => worker::console_warn!("{}: {error}", row.branch()),
295 }
296 }
297
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request298 /// Sends entries back to waiting, to be tested again.
299 async fn retest(&self, rows: &[&EntryRow]) -> Result<()> {
300 for row in rows {
301 self.db
302 .prepare(
303 "UPDATE queue_entries
304 SET state = 'waiting', token_hash = NULL, combined_commit = NULL,
305 results = NULL, error = NULL, ahead = NULL
306 WHERE id = ? AND state IN ('testing', 'passed')",
307 )
308 .bind(&[row.id.as_str().into()])?
309 .run()
310 .await?;
311 }
312 Ok(())
313 }
314
315 /// The next batch to test, if nothing is being tested now: one job per
316 /// entry, each building the default branch with that entry and every
317 /// entry ahead of it.
318 pub(crate) async fn queue_build(&self, a: QueueBuildArgs) -> Result<Vec<QueueJob>> {
319 let active = self.entries(&a.repo_id, true).await?;
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge320 // How the default branch's merge queue rule batches (rulesets.rs).
321 let rule = self.queue_rule(&a.repo_id).await?;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request322 // A batch that has taken too long is tested again.
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge323 let stale = minutes_ago(u64::from(rule.check_response_timeout_minutes).max(1));
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request324 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 }
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge338 let batch: Vec<EntryRow> = active.into_iter().take(rule.max_entries_to_build.max(1) as usize).collect();
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request339 let Some(first) = batch.first() else {
340 return Ok(Vec::new());
341 };
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge342 // Too few to start yet, and the oldest has not waited long enough.
343 if batch.len() < rule.min_entries_to_merge as usize
344 && first.created_at.as_str() > minutes_ago(u64::from(rule.min_entries_wait_minutes)).as_str()
345 {
346 return Ok(Vec::new());
347 }
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request348 let Some(actor) = first.actor() else {
349 return Ok(Vec::new());
350 };
351 let repo: Outcome<Repo> = g1t_kit::call(
352 &self.repos,
353 "get_by_id",
354 &GetByIdArgs {
355 id: a.repo_id.clone(),
356 viewer: Some(actor),
357 },
358 )
359 .await?;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look360 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request361 return Ok(Vec::new());
362 };
363 let base: Option<String> = g1t_kit::call(
364 &self.repos,
365 "head",
366 &HeadArgs {
367 repo_id: repo.id.clone(),
368 branch: repo.default_branch.clone(),
369 },
370 )
371 .await?;
372 let Some(base) = base else {
373 return Ok(Vec::new());
374 };
375
Fast pages, required checks on the branch, self-hosted runners, honest incidents376 // Each entry's change, as it is now.
377 let mut items: Vec<(EntryRow, Pull, QueueStackItem)> = Vec::new();
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request378 for row in batch {
379 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
380 continue;
381 };
382 if pull.status != PullStatus::Open {
383 self.leave(&repo.id, &pull, QueueState::Removed, Some("It was closed."))
384 .await?;
385 continue;
386 }
387 let head: Option<String> = g1t_kit::call(
388 &self.repos,
389 "head",
390 &HeadArgs {
391 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
392 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
393 },
394 )
395 .await?;
396 let Some(commit) = head else {
397 continue;
398 };
399 let source = pull.fork.clone().unwrap_or_else(|| RepoPath {
400 namespace: repo.namespace.clone(),
401 name: repo.name.clone(),
402 });
403 let item = QueueStackItem {
404 number: pull.number,
405 title: pull.title.clone(),
406 source,
407 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
408 commit,
409 };
Fast pages, required checks on the branch, self-hosted runners, honest incidents410 items.push((row, pull, item));
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request411 }
412
413 let now = rfc3339(now_ms());
414 let mut jobs = Vec::new();
415 for index in 0..items.len() {
Fast pages, required checks on the branch, self-hosted runners, honest incidents416 let (row, pull, _) = &items[index];
417 let stack: Vec<QueueStackItem> = items[..=index].iter().map(|(_, _, item)| item.clone()).collect();
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request418 let ahead: Vec<u32> = stack[..index].iter().map(|item| item.number).collect();
419 let token = new_token();
420 self.db
421 .prepare(
422 "UPDATE queue_entries
423 SET state = 'testing', token_hash = ?, head_commit = ?, base_commit = ?,
424 ahead = ?, combined_commit = NULL, results = NULL, error = NULL,
425 tested_at = ?
426 WHERE id = ? AND state = 'waiting'",
427 )
428 .bind(&[
429 hash(&token).into(),
430 stack[index].commit.as_str().into(),
431 base.as_str().into(),
432 serde_json::to_string(&ahead)?.into(),
433 now.as_str().into(),
434 row.id.as_str().into(),
435 ])?
436 .run()
437 .await?;
438 // 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 rights439 // the pull request's owner (whoever asked g1t for it, or its
440 // author): a real account, whose token carries its memberships.
441 // (Whoever queued it may be g1t itself.)
442 let actor = pull.owner().clone();
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request443 jobs.push(QueueJob {
444 entry_id: row.id.clone(),
445 token,
446 repo: RepoPath {
447 namespace: repo.namespace.clone(),
448 name: repo.name.clone(),
449 },
450 default_branch: repo.default_branch.clone(),
451 base_commit: base.clone(),
452 branch: row.branch(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents453 // The state is checked by its merge_group workflows.
454 contract_checks: Vec::new(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request455 stack,
Fast pages, required checks on the branch, self-hosted runners, honest incidents456 checks: Vec::new(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request457 actor,
458 });
459 }
460 Ok(jobs)
461 }
462
463 /// A sandbox's result for one combined state.
464 pub(crate) async fn report_queue(&self, a: ReportQueueArgs) -> Result<Outcome<QueueState>> {
465 let row = self
466 .db
467 .prepare("SELECT * FROM queue_entries WHERE id = ? AND token_hash = ?")
468 .bind(&[a.entry_id.as_str().into(), hash(&a.token).into()])?
469 .first::<EntryRow>(None)
470 .await?;
471 let Some(row) = row else {
472 return Ok(Outcome::fail(FailureCode::NotFound, "No such queue entry."));
473 };
474 if row.state() != QueueState::Testing {
475 return Ok(Outcome::fail(
476 FailureCode::Conflict,
477 "This state is no longer being tested.",
478 ));
479 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents480 let built = a.error.is_none() && a.results.iter().all(|result| result.passed);
481 // A state that was built runs the repository's `merge_group`
482 // workflows; it stays in testing until they finish (see `statuses`),
483 // and the branch's required checks must pass on it.
484 let workflows = match (&a.combined_commit, built) {
Sidebar: the panels really slide485 (Some(commit), true) => self.start_merge_group(&row, commit).await.unwrap_or(0),
486 _ => 0,
487 };
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge488 let required = self.settings_by_id(&row.repo_id).await?.required_checks;
Fast pages, required checks on the branch, self-hosted runners, honest incidents489 // Nothing runs on it, so the required checks never would report.
490 let unchecked = (built && workflows == 0 && !required.is_empty()).then(|| {
491 format!(
492 "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",
493 if required.len() == 1 { "check" } else { "checks" },
494 crate::statuses::list(&required)
495 )
496 });
497 let passed = built && unchecked.is_none();
498 let error = a.error.clone().or(unchecked);
Sidebar: the panels really slide499 let state = if !passed {
500 QueueState::Failed
501 } else if workflows > 0 {
502 QueueState::Testing
503 } else {
504 QueueState::Passed
505 };
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request506 self.db
507 .prepare(
508 "UPDATE queue_entries
509 SET state = ?, combined_commit = ?, results = ?, error = ?, token_hash = NULL,
510 finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END
511 WHERE id = ?",
512 )
513 .bind(&[
514 state.as_str().into(),
515 a.combined_commit.as_deref().map_or(JsValue::NULL, JsValue::from),
516 serde_json::to_string(&a.results)?.into(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents517 error.as_deref().map_or(JsValue::NULL, JsValue::from),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request518 state.as_str().into(),
519 rfc3339(now_ms()).into(),
520 row.id.as_str().into(),
521 ])?
522 .run()
523 .await?;
524 if !passed {
Fast pages, required checks on the branch, self-hosted runners, honest incidents525 let report = ReportQueueArgs { error, ..a };
526 self.eject(&row, &report).await?;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request527 }
528 self.settle(&row.repo_id).await?;
529 self.changed(&row.repo_id).await?;
530 Ok(Outcome::Ok(state))
531 }
532
Sidebar: the panels really slide533 /// Asks the actions service to run the repository's `merge_group`
534 /// workflows on a combined state. Returns how many runs started.
535 async fn start_merge_group(&self, row: &EntryRow, commit: &str) -> Result<u32> {
536 #[derive(serde::Deserialize)]
537 struct Started {
538 runs: u32,
539 }
540 let started: Outcome<Started> = g1t_kit::call(
541 &self.actions,
542 "merge_group",
543 &serde_json::json!({
544 "repoId": row.repo_id,
545 "entry": row.id,
546 "sha": commit,
547 "headRef": format!("refs/heads/{}", row.branch()),
548 "baseSha": row.base_commit,
549 "number": row.number,
550 "ahead": row.ahead(),
551 }),
552 )
553 .await?;
554 Ok(match started {
555 Outcome::Ok(started) => started.runs,
556 Outcome::Fail(_) => 0,
557 })
558 }
559
Fast pages, required checks on the branch, self-hosted runners, honest incidents560 /// Workflows on a combined state finished: it passes and lands in turn
561 /// when they all passed and so did every required check, or fails and
562 /// leaves the queue.
563 pub(crate) async fn merge_group_finished(
564 &self,
565 repo_id: &str,
566 commit: &str,
567 facts: &crate::statuses::WorkflowFacts,
568 ) -> Result<()> {
Sidebar: the panels really slide569 let row = self
570 .db
571 .prepare(
572 "SELECT * FROM queue_entries WHERE repo_id = ? AND combined_commit = ? AND state = 'testing' AND token_hash IS NULL",
573 )
574 .bind(&[repo_id.into(), commit.into()])?
575 .first::<EntryRow>(None)
576 .await?;
577 let Some(row) = row else { return Ok(()) };
Fast pages, required checks on the branch, self-hosted runners, honest incidents578 let missing = facts.expected();
579 let passed = facts.failed.is_empty() && facts.required_failed().is_empty() && missing.is_empty();
Sidebar: the panels really slide580 self.db
581 .prepare("UPDATE queue_entries SET state = ?, finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END WHERE id = ?")
582 .bind(&[
583 (if passed { "passed" } else { "failed" }).into(),
584 (if passed { "passed" } else { "failed" }).into(),
585 rfc3339(now_ms()).into(),
586 row.id.as_str().into(),
587 ])?
588 .run()
589 .await?;
590 if !passed {
591 let report = ReportQueueArgs {
592 entry_id: row.id.clone(),
593 token: String::new(),
594 combined_commit: Some(commit.to_owned()),
595 results: Vec::new(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents596 error: Some(if facts.failed.is_empty() {
597 format!(
598 "the required {} {} did not report on it. Add merge_group to the on: of the workflows the branch requires",
599 if missing.len() == 1 { "check" } else { "checks" },
600 crate::statuses::list(&missing)
601 )
602 } else {
603 format!("the workflow {} failed on it", crate::statuses::list(&facts.failed))
604 }),
Sidebar: the panels really slide605 conflict_with: None,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains606 conflicts: Vec::new(),
Sidebar: the panels really slide607 };
608 self.eject(&row, &report).await?;
609 }
610 self.settle(repo_id).await?;
611 self.changed(repo_id).await?;
612 Ok(())
613 }
614
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request615 /// An entry whose combined state failed: the entries tested on top of
616 /// it are tested again without it, and the failure is recorded as a
617 /// failed check run of its pull request, so a g1t agent is sent back.
618 async fn eject(&self, row: &EntryRow, report: &ReportQueueArgs) -> Result<()> {
Merge queue: tested states are deleted once their entry leaves619 self.drop_branch(row).await;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request620 let active = self.entries(&row.repo_id, true).await?;
621 let behind: Vec<&EntryRow> = active
622 .iter()
623 .filter(|other| other.state() != QueueState::Waiting && other.ahead().contains(&row.number))
624 .collect();
625 self.retest(&behind).await?;
626
627 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
628 return Ok(());
629 };
630 let ahead = row.ahead();
631 let state = if ahead.is_empty() {
632 "the default branch as it is now".to_owned()
633 } else {
634 format!(
635 "the default branch with {} merged in first",
636 ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ")
637 )
638 };
Agents and memory, checks and conflicts, profiles, slug renames, custom domains639 // Named in backticks, so the conversation can link each to the diff.
640 let files = crate::mergeability::tidy(report.conflicts.clone())
641 .iter()
642 .map(|path| format!("`{path}`"))
643 .collect::<Vec<_>>()
644 .join(", ");
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request645 let why = match (&report.error, report.conflict_with) {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains646 (_, Some(other)) if other != row.number && !files.is_empty() => format!(
647 "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."
648 ),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request649 (_, Some(other)) if other != row.number => format!(
650 "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."
651 ),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains652 (Some(_), None) if !files.is_empty() => format!(
653 "Its change conflicts with the default branch in {files}. Bring it up to date with the default branch, and merge it again."
654 ),
Sidebar: the panels really slide655 (Some(error), _) if error.starts_with("the workflow ") => {
656 format!("{} when it was combined with {state}.", error.replacen("the workflow", "The workflow", 1).trim_end_matches(" on it"))
657 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents658 (Some(error), _) if error.starts_with("the required ") => {
659 format!("Combined with {state}, {error}.")
660 }
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request661 (Some(error), _) => format!("Its combined state could not be built or checked: {error}"),
662 (None, _) => format!(
Fast pages, required checks on the branch, self-hosted runners, honest incidents663 "It failed when it was combined with {state}, though it may pass on its own."
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request664 ),
665 };
666 let now = now_ms();
667 let run_id = new_id("chk", now);
668 let mut results = report.results.clone();
669 for result in &mut results {
670 result.command = format!("{} (merge queue, on {state})", result.command);
671 }
672 self.db
673 .batch(vec![
674 self.db
675 .prepare(
676 "INSERT INTO check_runs
677 (id, pull_id, head_commit, status, results, error, token_hash,
678 created_at, finished_at)
679 VALUES (?, ?, ?, 'failed', ?, ?, ?, ?, ?)",
680 )
681 .bind(&[
682 run_id.as_str().into(),
683 pull.id.as_str().into(),
684 row.head_commit.as_deref().unwrap_or_default().into(),
685 serde_json::to_string(&results)?.into(),
686 why.as_str().into(),
687 hash(&new_token()).into(),
688 rfc3339(now).into(),
689 rfc3339(now).into(),
690 ])?,
691 self.db
692 .prepare("UPDATE pulls SET check_status = 'failed', check_run_id = ? WHERE id = ?")
693 .bind(&[run_id.as_str().into(), pull.id.as_str().into()])?,
694 ])
695 .await?;
696 self.note(
697 &row.repo_id,
698 pull.number,
699 ("g1t", "g1t"),
700 &format!("was taken out of the merge queue. {why}"),
701 )
702 .await?;
703 self.publish_as(
704 "checks.completed",
705 &row.repo_id,
706 None,
707 ChecksEvent {
708 pull_id: pull.id.clone(),
709 repo_id: row.repo_id.clone(),
710 number: pull.number,
711 status: "failed",
712 commit: row.head_commit.clone().unwrap_or_default(),
713 },
714 )
715 .await?;
716 Ok(())
717 }
718
719 /// Lands every entry at the front of the queue whose tested state
720 /// passed, in order.
721 async fn settle(&self, repo_id: &str) -> Result<()> {
722 let active = self.entries(repo_id, true).await?;
723 let Some(viewer) = active.first().and_then(EntryRow::actor) else {
724 return Ok(());
725 };
726 let repo: Outcome<Repo> = g1t_kit::call(
727 &self.repos,
728 "get_by_id",
729 &GetByIdArgs {
730 id: repo_id.to_owned(),
731 viewer: Some(viewer),
732 },
733 )
734 .await?;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look735 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request736 return Ok(());
737 };
738 for row in active {
739 if row.state() != QueueState::Passed {
740 // The front has not passed yet: nothing behind it may land.
741 return Ok(());
742 }
743 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
744 continue;
745 };
746 let Some(actor) = row.actor() else {
747 continue;
748 };
749 // Pushed to since it was tested: test it again as it is now.
750 if pull.status != PullStatus::Open {
751 self.leave(repo_id, &pull, QueueState::Removed, Some("It was closed."))
752 .await?;
753 continue;
754 }
755 let head: Option<String> = g1t_kit::call(
756 &self.repos,
757 "head",
758 &HeadArgs {
759 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo_id.to_owned()),
760 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
761 },
762 )
763 .await?;
764 if head.is_some() && head != row.head_commit {
765 let active = self.entries(repo_id, true).await?;
766 let again: Vec<&EntryRow> = active
767 .iter()
768 .filter(|other| other.id == row.id || other.ahead().contains(&row.number))
769 .collect();
770 self.retest(&again).await?;
771 return Ok(());
772 }
773 let landed: Outcome<Landed> = g1t_kit::call(
774 &self.repos,
775 "land",
776 &LandArgs {
777 source_id: repo_id.to_owned(),
778 branch: Some(row.branch()),
779 actor: actor.clone(),
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar780 // The queue lands on the default branch.
781 target_branch: None,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request782 },
783 )
784 .await?;
785 let landed = match landed {
786 Outcome::Ok(landed) => landed,
787 Outcome::Fail(_) => {
788 // The default branch moved outside the queue: every
789 // tested state is built on something that is gone.
790 let active = self.entries(repo_id, true).await?;
791 let all: Vec<&EntryRow> = active.iter().collect();
792 self.retest(&all).await?;
793 return Ok(());
794 }
795 };
796 self.db
797 .prepare(
798 "UPDATE queue_entries SET state = 'landed', finished_at = ? WHERE id = ?",
799 )
800 .bind(&[rfc3339(now_ms()).into(), row.id.as_str().into()])?
801 .run()
802 .await?;
Merge queue: tested states are deleted once their entry leaves803 self.drop_branch(&row).await;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request804 self.record_merge(&repo, pull, &actor, row.keep_issue_open != 0, landed)
805 .await?;
806 }
807 Ok(())
808 }
809}

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