g1t/services/work/src/queue.rs
//! The merge queue: pull requests are tested together with the ones ahead
//! of them, and only a combination that passed reaches the default branch.
//!
//! A batch of up to [`BATCH`] entries is tested at once, speculatively: one
//! sandbox per entry builds the default branch with that entry and every
//! entry ahead of it merged in, runs all of their acceptance checks, and
//! pushes the result to `g1t-queue/<entry>`. Entries land in order, each by
//! moving the default branch to its tested state, once everything ahead has
//! landed. One that fails leaves the queue with a failed check run, which
//! sends a g1t agent back to fix it; the entries behind it are tested again
//! without it.
use futures_util::future::try_join_all;
use g1t_contracts::events::{ChecksEvent, QueueChanged};
use g1t_contracts::repos::{DeleteBranchArgs, GetByIdArgs, HeadArgs, LandArgs, Landed, Repo, RepoPath};
use g1t_contracts::time::rfc3339;
use g1t_contracts::work::*;
use g1t_contracts::{FailureCode, Outcome, User, new_id};
use g1t_kit::now_ms;
use serde::Deserialize;
use worker::Result;
use worker::wasm_bindgen::JsValue;
use crate::Work;
use crate::checks::{hash, new_token};
/// How many entries are tested at once.
const BATCH: usize = 4;
/// How long a batch may take before it is tested again.
const TESTING_MINUTES: u64 = 45;
/// How many entries that left the queue are shown.
const RECENT: u32 = 20;
/// How many completed issues' checks make up the default branch's contract.
const CONTRACT_ISSUES: u32 = 30;
/// Where tested states are pushed, in the repository itself.
const BRANCH_PREFIX: &str = "g1t-queue/";
#[derive(Clone, Deserialize)]
pub(crate) struct EntryRow {
id: String,
repo_id: String,
pull_id: String,
number: u32,
state: String,
enqueued_by: String,
keep_issue_open: u32,
head_commit: Option<String>,
base_commit: Option<String>,
ahead: Option<String>,
combined_commit: Option<String>,
results: Option<String>,
error: Option<String>,
tested_at: Option<String>,
created_at: String,
finished_at: Option<String>,
}
impl EntryRow {
fn state(&self) -> QueueState {
match self.state.as_str() {
"testing" => QueueState::Testing,
"passed" => QueueState::Passed,
"failed" => QueueState::Failed,
"landed" => QueueState::Landed,
"removed" => QueueState::Removed,
_ => QueueState::Waiting,
}
}
fn actor(&self) -> Option<User> {
serde_json::from_str(&self.enqueued_by).ok()
}
fn ahead(&self) -> Vec<u32> {
self.ahead
.as_deref()
.and_then(|ahead| serde_json::from_str(ahead).ok())
.unwrap_or_default()
}
fn branch(&self) -> String {
format!("{BRANCH_PREFIX}{}", self.id)
}
}
fn minutes_ago(minutes: u64) -> String {
rfc3339(now_ms().saturating_sub(minutes * 60 * 1000))
}
impl Work {
async fn entries(&self, repo_id: &str, active: bool) -> Result<Vec<EntryRow>> {
let sql = if active {
"SELECT * FROM queue_entries
WHERE repo_id = ? AND state IN ('waiting', 'testing', 'passed')
ORDER BY created_at, id"
} else {
"SELECT * FROM queue_entries
WHERE repo_id = ? AND state IN ('failed', 'landed', 'removed')
ORDER BY finished_at DESC LIMIT ?"
};
let statement = self.db.prepare(sql);
let statement = if active {
statement.bind(&[repo_id.into()])?
} else {
statement.bind(&[repo_id.into(), RECENT.into()])?
};
statement.all().await?.results::<EntryRow>()
}
/// The entry a pull request has in the queue now, if any.
pub(crate) async fn queued_entry(&self, pull_id: &str) -> Result<Option<(QueueState, Vec<u32>)>> {
let row = self
.db
.prepare(
"SELECT * FROM queue_entries
WHERE pull_id = ? AND state IN ('waiting', 'testing', 'passed') LIMIT 1",
)
.bind(&[pull_id.into()])?
.first::<EntryRow>(None)
.await?;
Ok(row.map(|row| (row.state(), row.ahead())))
}
async fn changed(&self, repo_id: &str) -> Result<()> {
self.publish_as(
"queue.changed",
repo_id,
None,
QueueChanged {
repo_id: repo_id.to_owned(),
},
)
.await
}
/// Puts a pull request that may merge into the queue instead. Merging
/// it again while it is queued changes nothing.
pub(crate) async fn enqueue(
&self,
repo: &Repo,
pull: &Pull,
actor: &User,
keep_issue_open: bool,
) -> Result<Outcome<Pull>> {
if self.queued_entry(&pull.id).await?.is_some() {
return Ok(Outcome::Ok(pull.clone()));
}
let now = now_ms();
self.db
.prepare(
"INSERT INTO queue_entries
(id, repo_id, pull_id, number, state, enqueued_by, keep_issue_open, created_at)
VALUES (?, ?, ?, ?, 'waiting', ?, ?, ?)",
)
.bind(&[
new_id("qen", now).into(),
repo.id.as_str().into(),
pull.id.as_str().into(),
pull.number.into(),
serde_json::to_string(actor)?.into(),
u32::from(keep_issue_open).into(),
rfc3339(now).into(),
])?
.run()
.await?;
let who = (actor.id.as_str(), actor.username.as_str());
self.note(&repo.id, pull.number, who, "added this to the merge queue")
.await?;
self.changed(&repo.id).await?;
Ok(Outcome::Ok(pull.clone()))
}
pub(crate) async fn queue(&self, a: QueueArgs) -> Result<Outcome<QueueView>> {
let repo = match self.repo(&a.repo, &a.viewer).await? {
Outcome::Ok(repo) => repo,
Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
};
let enabled = self.settings(&repo.id).await?.merge_queue;
let (active, recent) = (self.entries(&repo.id, true).await?, self.entries(&repo.id, false).await?);
let numbers: Vec<u32> = active.iter().chain(&recent).map(|row| row.number).collect();
let pulls = try_join_all(numbers.iter().map(|number| self.pull(&repo.id, *number))).await?;
let view = |rows: Vec<EntryRow>, pulls: &[Option<Pull>]| -> Vec<QueueEntry> {
rows.into_iter()
.zip(pulls)
.map(|(row, pull)| QueueEntry {
state: row.state(),
ahead: row.ahead(),
enqueued_by: row.actor().map_or_else(|| "g1t".to_owned(), |actor| actor.username),
title: pull.as_ref().map(|p| p.title.clone()).unwrap_or_default(),
agent: pull.as_ref().map(|p| p.agent.clone()).unwrap_or_default(),
results: row
.results
.as_deref()
.and_then(|results| serde_json::from_str(results).ok())
.unwrap_or_default(),
id: row.id,
number: row.number,
base_commit: row.base_commit,
combined_commit: row.combined_commit,
error: row.error,
created_at: row.created_at,
finished_at: row.finished_at,
})
.collect()
};
let split = active.len();
Ok(Outcome::Ok(QueueView {
enabled,
active: view(active, &pulls[..split]),
recent: view(recent, &pulls[split..]),
}))
}
/// Takes a pull request out of the queue, by the hand of someone who
/// may merge.
pub(crate) async fn remove_from_queue(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
let viewer = Some(a.actor.clone());
let (repo, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
Outcome::Ok(found) => found,
Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
};
if let Outcome::Fail(failure) = crate::retired::writable(&repo) {
return Ok(Outcome::Fail(failure));
}
if !a.actor.verified {
return Ok(Outcome::fail(FailureCode::Forbidden, crate::UNVERIFIED));
}
if let Outcome::Fail(failure) =
crate::allowed(Some(&a.actor), &repo, g1t_contracts::access::Capability::Merge)
{
return Ok(Outcome::Fail(failure));
}
if self.leave(&repo.id, &pull, QueueState::Removed, None).await? {
let who = (a.actor.id.as_str(), a.actor.username.as_str());
self.note(&repo.id, pull.number, who, "removed this from the merge queue")
.await?;
self.changed(&repo.id).await?;
}
Ok(Outcome::Ok(pull))
}
/// Takes a pull request's entry out of the queue, and sends every entry
/// whose tested state included it back to waiting. Whether it had one.
pub(crate) async fn leave(
&self,
repo_id: &str,
pull: &Pull,
state: QueueState,
error: Option<&str>,
) -> Result<bool> {
let active = self.entries(repo_id, true).await?;
let Some(entry) = active.iter().find(|row| row.pull_id == pull.id) else {
return Ok(false);
};
let now = rfc3339(now_ms());
self.db
.prepare(
"UPDATE queue_entries SET state = ?, error = COALESCE(?, error), finished_at = ?
WHERE id = ?",
)
.bind(&[
state.as_str().into(),
error.map_or(JsValue::NULL, JsValue::from),
now.as_str().into(),
entry.id.as_str().into(),
])?
.run()
.await?;
self.drop_branch(entry).await;
let behind: Vec<&EntryRow> = active
.iter()
.filter(|row| row.state() != QueueState::Waiting && row.ahead().contains(&pull.number))
.collect();
self.retest(&behind).await?;
Ok(true)
}
/// Removes an entry's tested state from the repository once it has
/// left the queue, as GitHub does with its queue's branches. A branch
/// left behind is untidy, not wrong, so a failure only logs.
async fn drop_branch(&self, row: &EntryRow) {
let deleted: Result<Outcome<bool>> = g1t_kit::call(
&self.repos,
"delete_branch",
&DeleteBranchArgs {
repo_id: row.repo_id.clone(),
branch: row.branch(),
},
)
.await;
match deleted {
Ok(Outcome::Ok(_)) => {}
Ok(Outcome::Fail(failure)) => worker::console_warn!("{}: {}", row.branch(), failure.message),
Err(error) => worker::console_warn!("{}: {error}", row.branch()),
}
}
/// Sends entries back to waiting, to be tested again.
async fn retest(&self, rows: &[&EntryRow]) -> Result<()> {
for row in rows {
self.db
.prepare(
"UPDATE queue_entries
SET state = 'waiting', token_hash = NULL, combined_commit = NULL,
results = NULL, error = NULL, ahead = NULL
WHERE id = ? AND state IN ('testing', 'passed')",
)
.bind(&[row.id.as_str().into()])?
.run()
.await?;
}
Ok(())
}
/// The next batch to test, if nothing is being tested now: one job per
/// entry, each building the default branch with that entry and every
/// entry ahead of it.
pub(crate) async fn queue_build(&self, a: QueueBuildArgs) -> Result<Vec<QueueJob>> {
let active = self.entries(&a.repo_id, true).await?;
// A batch that has taken too long is tested again.
let stale = minutes_ago(TESTING_MINUTES);
let stuck: Vec<&EntryRow> = active
.iter()
.filter(|row| {
row.state() == QueueState::Testing
&& row.tested_at.as_deref().is_none_or(|at| at < stale.as_str())
})
.collect();
if !stuck.is_empty() {
self.retest(&stuck).await?;
return Box::pin(self.queue_build(a)).await;
}
if active.iter().any(|row| row.state() != QueueState::Waiting) {
return Ok(Vec::new());
}
let batch: Vec<EntryRow> = active.into_iter().take(BATCH).collect();
let Some(first) = batch.first() else {
return Ok(Vec::new());
};
let Some(actor) = first.actor() else {
return Ok(Vec::new());
};
let repo: Outcome<Repo> = g1t_kit::call(
&self.repos,
"get_by_id",
&GetByIdArgs {
id: a.repo_id.clone(),
viewer: Some(actor),
},
)
.await?;
let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
return Ok(Vec::new());
};
let base: Option<String> = g1t_kit::call(
&self.repos,
"head",
&HeadArgs {
repo_id: repo.id.clone(),
branch: repo.default_branch.clone(),
},
)
.await?;
let Some(base) = base else {
return Ok(Vec::new());
};
// Each entry's change, as it is now, and the checks it brings.
let mut items: Vec<(EntryRow, Pull, QueueStackItem, Vec<String>)> = Vec::new();
for row in batch {
let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
continue;
};
if pull.status != PullStatus::Open {
self.leave(&repo.id, &pull, QueueState::Removed, Some("It was closed."))
.await?;
continue;
}
let head: Option<String> = g1t_kit::call(
&self.repos,
"head",
&HeadArgs {
repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
},
)
.await?;
let Some(commit) = head else {
continue;
};
let checks = match pull.issue {
Some(number) => self
.issue(&repo.id, number)
.await?
.map(|issue| issue.checks)
.unwrap_or_default(),
None => Vec::new(),
};
let source = pull.fork.clone().unwrap_or_else(|| RepoPath {
namespace: repo.namespace.clone(),
name: repo.name.clone(),
});
let item = QueueStackItem {
number: pull.number,
title: pull.title.clone(),
source,
branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
commit,
};
items.push((row, pull, item, checks));
}
// What the default branch has promised so far: every check an issue
// passed when it landed, newest first.
#[derive(Deserialize)]
struct ChecksRow {
checks: String,
}
let mut contract: Vec<String> = Vec::new();
for row in self
.db
.prepare(
"SELECT checks FROM issues
WHERE repo_id = ? AND state = 'closed' AND reason = 'completed' AND checks != '[]'
ORDER BY closed_at DESC LIMIT ?",
)
.bind(&[repo.id.as_str().into(), CONTRACT_ISSUES.into()])?
.all()
.await?
.results::<ChecksRow>()?
{
for check in serde_json::from_str::<Vec<String>>(&row.checks).unwrap_or_default() {
if !contract.contains(&check) {
contract.push(check);
}
}
}
let now = rfc3339(now_ms());
let mut jobs = Vec::new();
for index in 0..items.len() {
let (row, pull, _, _) = &items[index];
let stack: Vec<QueueStackItem> = items[..=index].iter().map(|(_, _, item, _)| item.clone()).collect();
let mut checks: Vec<String> = Vec::new();
for (_, _, _, theirs) in &items[..=index] {
for check in theirs {
if !checks.contains(check) {
checks.push(check.clone());
}
}
}
let ahead: Vec<u32> = stack[..index].iter().map(|item| item.number).collect();
let token = new_token();
self.db
.prepare(
"UPDATE queue_entries
SET state = 'testing', token_hash = ?, head_commit = ?, base_commit = ?,
ahead = ?, combined_commit = NULL, results = NULL, error = NULL,
tested_at = ?
WHERE id = ? AND state = 'waiting'",
)
.bind(&[
hash(&token).into(),
stack[index].commit.as_str().into(),
base.as_str().into(),
serde_json::to_string(&ahead)?.into(),
now.as_str().into(),
row.id.as_str().into(),
])?
.run()
.await?;
// The sandbox reads the change and pushes the tested state as
// the pull request's author: a real account, whose token carries
// its memberships. (Whoever queued it may be g1t itself.)
let actor = pull.author.clone();
jobs.push(QueueJob {
entry_id: row.id.clone(),
token,
repo: RepoPath {
namespace: repo.namespace.clone(),
name: repo.name.clone(),
},
default_branch: repo.default_branch.clone(),
base_commit: base.clone(),
branch: row.branch(),
contract_checks: contract.iter().filter(|check| !checks.contains(check)).cloned().collect(),
stack,
checks,
actor,
});
}
Ok(jobs)
}
/// A sandbox's result for one combined state.
pub(crate) async fn report_queue(&self, a: ReportQueueArgs) -> Result<Outcome<QueueState>> {
let row = self
.db
.prepare("SELECT * FROM queue_entries WHERE id = ? AND token_hash = ?")
.bind(&[a.entry_id.as_str().into(), hash(&a.token).into()])?
.first::<EntryRow>(None)
.await?;
let Some(row) = row else {
return Ok(Outcome::fail(FailureCode::NotFound, "No such queue entry."));
};
if row.state() != QueueState::Testing {
return Ok(Outcome::fail(
FailureCode::Conflict,
"This state is no longer being tested.",
));
}
let passed = a.error.is_none() && a.results.iter().all(|result| result.passed);
// A state that passed its checks still runs the repository's
// `merge_group` workflows, as GitHub's merge queue does; it stays in
// testing until they finish (see `statuses`).
let workflows = match (&a.combined_commit, passed) {
(Some(commit), true) => self.start_merge_group(&row, commit).await.unwrap_or(0),
_ => 0,
};
let state = if !passed {
QueueState::Failed
} else if workflows > 0 {
QueueState::Testing
} else {
QueueState::Passed
};
self.db
.prepare(
"UPDATE queue_entries
SET state = ?, combined_commit = ?, results = ?, error = ?, token_hash = NULL,
finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END
WHERE id = ?",
)
.bind(&[
state.as_str().into(),
a.combined_commit.as_deref().map_or(JsValue::NULL, JsValue::from),
serde_json::to_string(&a.results)?.into(),
a.error.as_deref().map_or(JsValue::NULL, JsValue::from),
state.as_str().into(),
rfc3339(now_ms()).into(),
row.id.as_str().into(),
])?
.run()
.await?;
if !passed {
self.eject(&row, &a).await?;
}
self.settle(&row.repo_id).await?;
self.changed(&row.repo_id).await?;
Ok(Outcome::Ok(state))
}
/// Asks the actions service to run the repository's `merge_group`
/// workflows on a combined state. Returns how many runs started.
async fn start_merge_group(&self, row: &EntryRow, commit: &str) -> Result<u32> {
#[derive(serde::Deserialize)]
struct Started {
runs: u32,
}
let started: Outcome<Started> = g1t_kit::call(
&self.actions,
"merge_group",
&serde_json::json!({
"repoId": row.repo_id,
"entry": row.id,
"sha": commit,
"headRef": format!("refs/heads/{}", row.branch()),
"baseSha": row.base_commit,
"number": row.number,
"ahead": row.ahead(),
}),
)
.await?;
Ok(match started {
Outcome::Ok(started) => started.runs,
Outcome::Fail(_) => 0,
})
}
/// Workflows on a combined state finished: it passes and lands in turn,
/// or fails and leaves the queue as for failed checks.
pub(crate) async fn merge_group_finished(&self, repo_id: &str, commit: &str, failed: &[String]) -> Result<()> {
let row = self
.db
.prepare(
"SELECT * FROM queue_entries WHERE repo_id = ? AND combined_commit = ? AND state = 'testing' AND token_hash IS NULL",
)
.bind(&[repo_id.into(), commit.into()])?
.first::<EntryRow>(None)
.await?;
let Some(row) = row else { return Ok(()) };
let passed = failed.is_empty();
self.db
.prepare("UPDATE queue_entries SET state = ?, finished_at = CASE WHEN ? = 'failed' THEN ? ELSE NULL END WHERE id = ?")
.bind(&[
(if passed { "passed" } else { "failed" }).into(),
(if passed { "passed" } else { "failed" }).into(),
rfc3339(now_ms()).into(),
row.id.as_str().into(),
])?
.run()
.await?;
if !passed {
let report = ReportQueueArgs {
entry_id: row.id.clone(),
token: String::new(),
combined_commit: Some(commit.to_owned()),
results: Vec::new(),
error: Some(format!(
"the workflow {} failed on it",
crate::statuses::list(failed)
)),
conflict_with: None,
conflicts: Vec::new(),
};
self.eject(&row, &report).await?;
}
self.settle(repo_id).await?;
self.changed(repo_id).await?;
Ok(())
}
/// An entry whose combined state failed: the entries tested on top of
/// it are tested again without it, and the failure is recorded as a
/// failed check run of its pull request, so a g1t agent is sent back.
async fn eject(&self, row: &EntryRow, report: &ReportQueueArgs) -> Result<()> {
self.drop_branch(row).await;
let active = self.entries(&row.repo_id, true).await?;
let behind: Vec<&EntryRow> = active
.iter()
.filter(|other| other.state() != QueueState::Waiting && other.ahead().contains(&row.number))
.collect();
self.retest(&behind).await?;
let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
return Ok(());
};
let ahead = row.ahead();
let state = if ahead.is_empty() {
"the default branch as it is now".to_owned()
} else {
format!(
"the default branch with {} merged in first",
ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", ")
)
};
// Named in backticks, so the conversation can link each to the diff.
let files = crate::mergeability::tidy(report.conflicts.clone())
.iter()
.map(|path| format!("`{path}`"))
.collect::<Vec<_>>()
.join(", ");
let why = match (&report.error, report.conflict_with) {
(_, Some(other)) if other != row.number && !files.is_empty() => format!(
"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."
),
(_, Some(other)) if other != row.number => format!(
"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."
),
(Some(_), None) if !files.is_empty() => format!(
"Its change conflicts with the default branch in {files}. Bring it up to date with the default branch, and merge it again."
),
(Some(error), _) if error.starts_with("the workflow ") => {
format!("{} when it was combined with {state}.", error.replacen("the workflow", "The workflow", 1).trim_end_matches(" on it"))
}
(Some(error), _) => format!("Its combined state could not be built or checked: {error}"),
(None, _) => format!(
"The acceptance checks failed when it was combined with {state}, though it may pass on its own."
),
};
let now = now_ms();
let run_id = new_id("chk", now);
let mut results = report.results.clone();
for result in &mut results {
result.command = format!("{} (merge queue, on {state})", result.command);
}
self.db
.batch(vec![
self.db
.prepare(
"INSERT INTO check_runs
(id, pull_id, head_commit, status, results, error, token_hash,
created_at, finished_at)
VALUES (?, ?, ?, 'failed', ?, ?, ?, ?, ?)",
)
.bind(&[
run_id.as_str().into(),
pull.id.as_str().into(),
row.head_commit.as_deref().unwrap_or_default().into(),
serde_json::to_string(&results)?.into(),
why.as_str().into(),
hash(&new_token()).into(),
rfc3339(now).into(),
rfc3339(now).into(),
])?,
self.db
.prepare("UPDATE pulls SET check_status = 'failed', check_run_id = ? WHERE id = ?")
.bind(&[run_id.as_str().into(), pull.id.as_str().into()])?,
])
.await?;
self.note(
&row.repo_id,
pull.number,
("g1t", "g1t"),
&format!("was taken out of the merge queue. {why}"),
)
.await?;
self.publish_as(
"checks.completed",
&row.repo_id,
None,
ChecksEvent {
pull_id: pull.id.clone(),
repo_id: row.repo_id.clone(),
number: pull.number,
status: "failed",
commit: row.head_commit.clone().unwrap_or_default(),
},
)
.await?;
Ok(())
}
/// Lands every entry at the front of the queue whose tested state
/// passed, in order.
async fn settle(&self, repo_id: &str) -> Result<()> {
let active = self.entries(repo_id, true).await?;
let Some(viewer) = active.first().and_then(EntryRow::actor) else {
return Ok(());
};
let repo: Outcome<Repo> = g1t_kit::call(
&self.repos,
"get_by_id",
&GetByIdArgs {
id: repo_id.to_owned(),
viewer: Some(viewer),
},
)
.await?;
let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
return Ok(());
};
for row in active {
if row.state() != QueueState::Passed {
// The front has not passed yet: nothing behind it may land.
return Ok(());
}
let Some(pull) = self.pull_by_id(&row.pull_id).await? else {
continue;
};
let Some(actor) = row.actor() else {
continue;
};
// Pushed to since it was tested: test it again as it is now.
if pull.status != PullStatus::Open {
self.leave(repo_id, &pull, QueueState::Removed, Some("It was closed."))
.await?;
continue;
}
let head: Option<String> = g1t_kit::call(
&self.repos,
"head",
&HeadArgs {
repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo_id.to_owned()),
branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
},
)
.await?;
if head.is_some() && head != row.head_commit {
let active = self.entries(repo_id, true).await?;
let again: Vec<&EntryRow> = active
.iter()
.filter(|other| other.id == row.id || other.ahead().contains(&row.number))
.collect();
self.retest(&again).await?;
return Ok(());
}
let landed: Outcome<Landed> = g1t_kit::call(
&self.repos,
"land",
&LandArgs {
source_id: repo_id.to_owned(),
branch: Some(row.branch()),
actor: actor.clone(),
},
)
.await?;
let landed = match landed {
Outcome::Ok(landed) => landed,
Outcome::Fail(_) => {
// The default branch moved outside the queue: every
// tested state is built on something that is gone.
let active = self.entries(repo_id, true).await?;
let all: Vec<&EntryRow> = active.iter().collect();
self.retest(&all).await?;
return Ok(());
}
};
self.db
.prepare(
"UPDATE queue_entries SET state = 'landed', finished_at = ? WHERE id = ?",
)
.bind(&[rfc3339(now_ms()).into(), row.id.as_str().into()])?
.run()
.await?;
self.drop_branch(&row).await;
self.record_merge(&repo, pull, &actor, row.keep_issue_open != 0, landed)
.await?;
}
Ok(())
}
}