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

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