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

797 lines31,747 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//!
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
Fast pages, required checks on the branch, self-hosted runners, honest incidents6//! 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.
Agents as a team: lifecycle, merge queue, billing and a new shell13
14use futures_util::future::try_join_all;
15use g1t_contracts::events::{ChecksEvent, QueueChanged};
Merge queue: tested states are deleted once their entry leaves16use g1t_contracts::repos::{DeleteBranchArgs, GetByIdArgs, HeadArgs, LandArgs, Landed, Repo, RepoPath};
Agents as a team: lifecycle, merge queue, billing and a new shell17use 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)]
Sidebar: the panels really slide38pub(crate) struct EntryRow {
Agents as a team: lifecycle, merge queue, billing and a new shell39 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 = self
112 .db
113 .prepare(
114 "SELECT * FROM queue_entries
115 WHERE pull_id = ? AND state IN ('waiting', 'testing', 'passed') LIMIT 1",
116 )
117 .bind(&[pull_id.into()])?
118 .first::<EntryRow>(None)
119 .await?;
120 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 };
177 let enabled = self.settings(&repo.id).await?.merge_queue;
178 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?;
319 // A batch that has taken too long is tested again.
320 let stale = minutes_ago(TESTING_MINUTES);
321 let stuck: Vec<&EntryRow> = active
322 .iter()
323 .filter(|row| {
324 row.state() == QueueState::Testing
325 && row.tested_at.as_deref().is_none_or(|at| at < stale.as_str())
326 })
327 .collect();
328 if !stuck.is_empty() {
329 self.retest(&stuck).await?;
330 return Box::pin(self.queue_build(a)).await;
331 }
332 if active.iter().any(|row| row.state() != QueueState::Waiting) {
333 return Ok(Vec::new());
334 }
335 let batch: Vec<EntryRow> = active.into_iter().take(BATCH).collect();
336 let Some(first) = batch.first() else {
337 return Ok(Vec::new());
338 };
339 let Some(actor) = first.actor() else {
340 return Ok(Vec::new());
341 };
342 let repo: Outcome<Repo> = g1t_kit::call(
343 &self.repos,
344 "get_by_id",
345 &GetByIdArgs {
346 id: a.repo_id.clone(),
347 viewer: Some(actor),
348 },
349 )
350 .await?;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look351 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
Agents as a team: lifecycle, merge queue, billing and a new shell352 return Ok(Vec::new());
353 };
354 let base: Option<String> = g1t_kit::call(
355 &self.repos,
356 "head",
357 &HeadArgs {
358 repo_id: repo.id.clone(),
359 branch: repo.default_branch.clone(),
360 },
361 )
362 .await?;
363 let Some(base) = base else {
364 return Ok(Vec::new());
365 };
366
Fast pages, required checks on the branch, self-hosted runners, honest incidents367 // Each entry's change, as it is now.
368 let mut items: Vec<(EntryRow, Pull, QueueStackItem)> = Vec::new();
Agents as a team: lifecycle, merge queue, billing and a new shell369 for row in batch {
370 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
371 continue;
372 };
373 if pull.status != PullStatus::Open {
374 self.leave(&repo.id, &pull, QueueState::Removed, Some("It was closed."))
375 .await?;
376 continue;
377 }
378 let head: Option<String> = g1t_kit::call(
379 &self.repos,
380 "head",
381 &HeadArgs {
382 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
383 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
384 },
385 )
386 .await?;
387 let Some(commit) = head else {
388 continue;
389 };
390 let source = pull.fork.clone().unwrap_or_else(|| RepoPath {
391 namespace: repo.namespace.clone(),
392 name: repo.name.clone(),
393 });
394 let item = QueueStackItem {
395 number: pull.number,
396 title: pull.title.clone(),
397 source,
398 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
399 commit,
400 };
Fast pages, required checks on the branch, self-hosted runners, honest incidents401 items.push((row, pull, item));
Agents as a team: lifecycle, merge queue, billing and a new shell402 }
403
404 let now = rfc3339(now_ms());
405 let mut jobs = Vec::new();
406 for index in 0..items.len() {
Fast pages, required checks on the branch, self-hosted runners, honest incidents407 let (row, pull, _) = &items[index];
408 let stack: Vec<QueueStackItem> = items[..=index].iter().map(|(_, _, item)| item.clone()).collect();
Agents as a team: lifecycle, merge queue, billing and a new shell409 let ahead: Vec<u32> = stack[..index].iter().map(|item| item.number).collect();
410 let token = new_token();
411 self.db
412 .prepare(
413 "UPDATE queue_entries
414 SET state = 'testing', token_hash = ?, head_commit = ?, base_commit = ?,
415 ahead = ?, combined_commit = NULL, results = NULL, error = NULL,
416 tested_at = ?
417 WHERE id = ? AND state = 'waiting'",
418 )
419 .bind(&[
420 hash(&token).into(),
421 stack[index].commit.as_str().into(),
422 base.as_str().into(),
423 serde_json::to_string(&ahead)?.into(),
424 now.as_str().into(),
425 row.id.as_str().into(),
426 ])?
427 .run()
428 .await?;
429 // The sandbox reads the change and pushes the tested state as
430 // the pull request's author: a real account, whose token carries
431 // its memberships. (Whoever queued it may be g1t itself.)
432 let actor = pull.author.clone();
433 jobs.push(QueueJob {
434 entry_id: row.id.clone(),
435 token,
436 repo: RepoPath {
437 namespace: repo.namespace.clone(),
438 name: repo.name.clone(),
439 },
440 default_branch: repo.default_branch.clone(),
441 base_commit: base.clone(),
442 branch: row.branch(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents443 // The state is checked by its merge_group workflows.
444 contract_checks: Vec::new(),
Agents as a team: lifecycle, merge queue, billing and a new shell445 stack,
Fast pages, required checks on the branch, self-hosted runners, honest incidents446 checks: Vec::new(),
Agents as a team: lifecycle, merge queue, billing and a new shell447 actor,
448 });
449 }
450 Ok(jobs)
451 }
452
453 /// A sandbox's result for one combined state.
454 pub(crate) async fn report_queue(&self, a: ReportQueueArgs) -> Result<Outcome<QueueState>> {
455 let row = self
456 .db
457 .prepare("SELECT * FROM queue_entries WHERE id = ? AND token_hash = ?")
458 .bind(&[a.entry_id.as_str().into(), hash(&a.token).into()])?
459 .first::<EntryRow>(None)
460 .await?;
461 let Some(row) = row else {
462 return Ok(Outcome::fail(FailureCode::NotFound, "No such queue entry."));
463 };
464 if row.state() != QueueState::Testing {
465 return Ok(Outcome::fail(
466 FailureCode::Conflict,
467 "This state is no longer being tested.",
468 ));
469 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents470 let built = a.error.is_none() && a.results.iter().all(|result| result.passed);
471 // A state that was built runs the repository's `merge_group`
472 // workflows; it stays in testing until they finish (see `statuses`),
473 // and the branch's required checks must pass on it.
474 let workflows = match (&a.combined_commit, built) {
Sidebar: the panels really slide475 (Some(commit), true) => self.start_merge_group(&row, commit).await.unwrap_or(0),
476 _ => 0,
477 };
Fast pages, required checks on the branch, self-hosted runners, honest incidents478 let required = self.settings(&row.repo_id).await?.required_checks;
479 // Nothing runs on it, so the required checks never would report.
480 let unchecked = (built && workflows == 0 && !required.is_empty()).then(|| {
481 format!(
482 "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",
483 if required.len() == 1 { "check" } else { "checks" },
484 crate::statuses::list(&required)
485 )
486 });
487 let passed = built && unchecked.is_none();
488 let error = a.error.clone().or(unchecked);
Sidebar: the panels really slide489 let state = if !passed {
490 QueueState::Failed
491 } else if workflows > 0 {
492 QueueState::Testing
493 } else {
494 QueueState::Passed
495 };
Agents as a team: lifecycle, merge queue, billing and a new shell496 self.db
497 .prepare(
498 "UPDATE queue_entries
499 SET state = ?, combined_commit = ?, results = ?, error = ?, token_hash = NULL,
500 finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END
501 WHERE id = ?",
502 )
503 .bind(&[
504 state.as_str().into(),
505 a.combined_commit.as_deref().map_or(JsValue::NULL, JsValue::from),
506 serde_json::to_string(&a.results)?.into(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents507 error.as_deref().map_or(JsValue::NULL, JsValue::from),
Agents as a team: lifecycle, merge queue, billing and a new shell508 state.as_str().into(),
509 rfc3339(now_ms()).into(),
510 row.id.as_str().into(),
511 ])?
512 .run()
513 .await?;
514 if !passed {
Fast pages, required checks on the branch, self-hosted runners, honest incidents515 let report = ReportQueueArgs { error, ..a };
516 self.eject(&row, &report).await?;
Agents as a team: lifecycle, merge queue, billing and a new shell517 }
518 self.settle(&row.repo_id).await?;
519 self.changed(&row.repo_id).await?;
520 Ok(Outcome::Ok(state))
521 }
522
Sidebar: the panels really slide523 /// Asks the actions service to run the repository's `merge_group`
524 /// workflows on a combined state. Returns how many runs started.
525 async fn start_merge_group(&self, row: &EntryRow, commit: &str) -> Result<u32> {
526 #[derive(serde::Deserialize)]
527 struct Started {
528 runs: u32,
529 }
530 let started: Outcome<Started> = g1t_kit::call(
531 &self.actions,
532 "merge_group",
533 &serde_json::json!({
534 "repoId": row.repo_id,
535 "entry": row.id,
536 "sha": commit,
537 "headRef": format!("refs/heads/{}", row.branch()),
538 "baseSha": row.base_commit,
539 "number": row.number,
540 "ahead": row.ahead(),
541 }),
542 )
543 .await?;
544 Ok(match started {
545 Outcome::Ok(started) => started.runs,
546 Outcome::Fail(_) => 0,
547 })
548 }
549
Fast pages, required checks on the branch, self-hosted runners, honest incidents550 /// Workflows on a combined state finished: it passes and lands in turn
551 /// when they all passed and so did every required check, or fails and
552 /// leaves the queue.
553 pub(crate) async fn merge_group_finished(
554 &self,
555 repo_id: &str,
556 commit: &str,
557 facts: &crate::statuses::WorkflowFacts,
558 ) -> Result<()> {
Sidebar: the panels really slide559 let row = self
560 .db
561 .prepare(
562 "SELECT * FROM queue_entries WHERE repo_id = ? AND combined_commit = ? AND state = 'testing' AND token_hash IS NULL",
563 )
564 .bind(&[repo_id.into(), commit.into()])?
565 .first::<EntryRow>(None)
566 .await?;
567 let Some(row) = row else { return Ok(()) };
Fast pages, required checks on the branch, self-hosted runners, honest incidents568 let missing = facts.expected();
569 let passed = facts.failed.is_empty() && facts.required_failed().is_empty() && missing.is_empty();
Sidebar: the panels really slide570 self.db
571 .prepare("UPDATE queue_entries SET state = ?, finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END WHERE id = ?")
572 .bind(&[
573 (if passed { "passed" } else { "failed" }).into(),
574 (if passed { "passed" } else { "failed" }).into(),
575 rfc3339(now_ms()).into(),
576 row.id.as_str().into(),
577 ])?
578 .run()
579 .await?;
580 if !passed {
581 let report = ReportQueueArgs {
582 entry_id: row.id.clone(),
583 token: String::new(),
584 combined_commit: Some(commit.to_owned()),
585 results: Vec::new(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents586 error: Some(if facts.failed.is_empty() {
587 format!(
588 "the required {} {} did not report on it. Add merge_group to the on: of the workflows the branch requires",
589 if missing.len() == 1 { "check" } else { "checks" },
590 crate::statuses::list(&missing)
591 )
592 } else {
593 format!("the workflow {} failed on it", crate::statuses::list(&facts.failed))
594 }),
Sidebar: the panels really slide595 conflict_with: None,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains596 conflicts: Vec::new(),
Sidebar: the panels really slide597 };
598 self.eject(&row, &report).await?;
599 }
600 self.settle(repo_id).await?;
601 self.changed(repo_id).await?;
602 Ok(())
603 }
604
Agents as a team: lifecycle, merge queue, billing and a new shell605 /// An entry whose combined state failed: the entries tested on top of
606 /// it are tested again without it, and the failure is recorded as a
607 /// failed check run of its pull request, so a g1t agent is sent back.
608 async fn eject(&self, row: &EntryRow, report: &ReportQueueArgs) -> Result<()> {
Merge queue: tested states are deleted once their entry leaves609 self.drop_branch(row).await;
Agents as a team: lifecycle, merge queue, billing and a new shell610 let active = self.entries(&row.repo_id, true).await?;
611 let behind: Vec<&EntryRow> = active
612 .iter()
613 .filter(|other| other.state() != QueueState::Waiting && other.ahead().contains(&row.number))
614 .collect();
615 self.retest(&behind).await?;
616
617 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
618 return Ok(());
619 };
620 let ahead = row.ahead();
621 let state = if ahead.is_empty() {
622 "the default branch as it is now".to_owned()
623 } else {
624 format!(
625 "the default branch with {} merged in first",
626 ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ")
627 )
628 };
Agents and memory, checks and conflicts, profiles, slug renames, custom domains629 // Named in backticks, so the conversation can link each to the diff.
630 let files = crate::mergeability::tidy(report.conflicts.clone())
631 .iter()
632 .map(|path| format!("`{path}`"))
633 .collect::<Vec<_>>()
634 .join(", ");
Agents as a team: lifecycle, merge queue, billing and a new shell635 let why = match (&report.error, report.conflict_with) {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains636 (_, Some(other)) if other != row.number && !files.is_empty() => format!(
637 "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."
638 ),
Agents as a team: lifecycle, merge queue, billing and a new shell639 (_, Some(other)) if other != row.number => format!(
640 "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."
641 ),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains642 (Some(_), None) if !files.is_empty() => format!(
643 "Its change conflicts with the default branch in {files}. Bring it up to date with the default branch, and merge it again."
644 ),
Sidebar: the panels really slide645 (Some(error), _) if error.starts_with("the workflow ") => {
646 format!("{} when it was combined with {state}.", error.replacen("the workflow", "The workflow", 1).trim_end_matches(" on it"))
647 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents648 (Some(error), _) if error.starts_with("the required ") => {
649 format!("Combined with {state}, {error}.")
650 }
Agents as a team: lifecycle, merge queue, billing and a new shell651 (Some(error), _) => format!("Its combined state could not be built or checked: {error}"),
652 (None, _) => format!(
Fast pages, required checks on the branch, self-hosted runners, honest incidents653 "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 shell654 ),
655 };
656 let now = now_ms();
657 let run_id = new_id("chk", now);
658 let mut results = report.results.clone();
659 for result in &mut results {
660 result.command = format!("{} (merge queue, on {state})", result.command);
661 }
662 self.db
663 .batch(vec![
664 self.db
665 .prepare(
666 "INSERT INTO check_runs
667 (id, pull_id, head_commit, status, results, error, token_hash,
668 created_at, finished_at)
669 VALUES (?, ?, ?, 'failed', ?, ?, ?, ?, ?)",
670 )
671 .bind(&[
672 run_id.as_str().into(),
673 pull.id.as_str().into(),
674 row.head_commit.as_deref().unwrap_or_default().into(),
675 serde_json::to_string(&results)?.into(),
676 why.as_str().into(),
677 hash(&new_token()).into(),
678 rfc3339(now).into(),
679 rfc3339(now).into(),
680 ])?,
681 self.db
682 .prepare("UPDATE pulls SET check_status = 'failed', check_run_id = ? WHERE id = ?")
683 .bind(&[run_id.as_str().into(), pull.id.as_str().into()])?,
684 ])
685 .await?;
686 self.note(
687 &row.repo_id,
688 pull.number,
689 ("g1t", "g1t"),
690 &format!("was taken out of the merge queue. {why}"),
691 )
692 .await?;
693 self.publish_as(
694 "checks.completed",
695 &row.repo_id,
696 None,
697 ChecksEvent {
698 pull_id: pull.id.clone(),
699 repo_id: row.repo_id.clone(),
700 number: pull.number,
701 status: "failed",
702 commit: row.head_commit.clone().unwrap_or_default(),
703 },
704 )
705 .await?;
706 Ok(())
707 }
708
709 /// Lands every entry at the front of the queue whose tested state
710 /// passed, in order.
711 async fn settle(&self, repo_id: &str) -> Result<()> {
712 let active = self.entries(repo_id, true).await?;
713 let Some(viewer) = active.first().and_then(EntryRow::actor) else {
714 return Ok(());
715 };
716 let repo: Outcome<Repo> = g1t_kit::call(
717 &self.repos,
718 "get_by_id",
719 &GetByIdArgs {
720 id: repo_id.to_owned(),
721 viewer: Some(viewer),
722 },
723 )
724 .await?;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look725 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
Agents as a team: lifecycle, merge queue, billing and a new shell726 return Ok(());
727 };
728 for row in active {
729 if row.state() != QueueState::Passed {
730 // The front has not passed yet: nothing behind it may land.
731 return Ok(());
732 }
733 let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
734 continue;
735 };
736 let Some(actor) = row.actor() else {
737 continue;
738 };
739 // Pushed to since it was tested: test it again as it is now.
740 if pull.status != PullStatus::Open {
741 self.leave(repo_id, &pull, QueueState::Removed, Some("It was closed."))
742 .await?;
743 continue;
744 }
745 let head: Option<String> = g1t_kit::call(
746 &self.repos,
747 "head",
748 &HeadArgs {
749 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo_id.to_owned()),
750 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
751 },
752 )
753 .await?;
754 if head.is_some() && head != row.head_commit {
755 let active = self.entries(repo_id, true).await?;
756 let again: Vec<&EntryRow> = active
757 .iter()
758 .filter(|other| other.id == row.id || other.ahead().contains(&row.number))
759 .collect();
760 self.retest(&again).await?;
761 return Ok(());
762 }
763 let landed: Outcome<Landed> = g1t_kit::call(
764 &self.repos,
765 "land",
766 &LandArgs {
767 source_id: repo_id.to_owned(),
768 branch: Some(row.branch()),
769 actor: actor.clone(),
770 },
771 )
772 .await?;
773 let landed = match landed {
774 Outcome::Ok(landed) => landed,
775 Outcome::Fail(_) => {
776 // The default branch moved outside the queue: every
777 // tested state is built on something that is gone.
778 let active = self.entries(repo_id, true).await?;
779 let all: Vec<&EntryRow> = active.iter().collect();
780 self.retest(&all).await?;
781 return Ok(());
782 }
783 };
784 self.db
785 .prepare(
786 "UPDATE queue_entries SET state = 'landed', finished_at = ? WHERE id = ?",
787 )
788 .bind(&[rfc3339(now_ms()).into(), row.id.as_str().into()])?
789 .run()
790 .await?;
Merge queue: tested states are deleted once their entry leaves791 self.drop_branch(&row).await;
Agents as a team: lifecycle, merge queue, billing and a new shell792 self.record_merge(&repo, pull, &actor, row.keep_issue_open != 0, landed)
793 .await?;
794 }
795 Ok(())
796 }
797}