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/mergeability.rs

545 lines20,692 bytesCodeBlame
1//! Whether a pull request merges cleanly into the branch it targets, worked
2//! out ahead of time, as GitHub does, rather than discovered when a merge
3//! or the merge queue trips over it.
4//!
5//! It is worked out whenever the pull request's head or its target branch
6//! moves, once per pair of commits. The first pass needs no sandbox: if the
7//! files the pull request changed since it and the branch last agreed share
8//! none with the files the branch changed since then, the merge cannot
9//! conflict. Otherwise a `pull.mergecheck` event asks the runner for a short
10//! probe in a sandbox, which merges the two without an agent and reports the
11//! files that conflict. Only so many probes run at once per repository; the
12//! rest wait their turn and start as earlier ones report.
13
14use std::collections::HashSet;
15
16use g1t_contracts::repos::{BehindArgs, Divergence, GetByIdArgs, Repo, RepoPath};
17use g1t_contracts::time::rfc3339;
18use g1t_contracts::work::*;
19use g1t_contracts::{FailureCode, Outcome};
20use g1t_kit::now_ms;
21use serde::Deserialize;
22use worker::Result;
23
24use crate::Work;
25use crate::checks::{hash, new_token};
26use crate::rows::{NumberRow, ValueRow};
27
28/// How long a probe is waited for before it is asked for again.
29const PROBE_MINUTES: u64 = 10;
30/// How many probes one repository may have running at once.
31const MAX_PROBES: u32 = 3;
32/// How many conflicting files are kept.
33const MAX_CONFLICTS: usize = 100;
34/// How many open pull requests are looked at when their target moves.
35const MAX_TARGETING: u32 = 100;
36
37#[derive(Deserialize)]
38struct MergeRow {
39 mergeable: Option<String>,
40 conflicts: Option<String>,
41 mergeable_key: Option<String>,
42 mergeable_until: Option<String>,
43}
44
45#[derive(Deserialize)]
46struct ProbeRow {
47 id: String,
48 repo_id: String,
49 mergeable_key: Option<String>,
50 mergeable_token_hash: Option<String>,
51}
52
53/// The pair of commits an answer is for.
54fn key(divergence: &Divergence) -> String {
55 format!("{}..{}", divergence.head, divergence.base)
56}
57
58/// Whether the cheap first pass leaves the question open, so that a probe
59/// has to merge the two. It is settled when the branch has not moved, or
60/// when the two sides changed no file in common.
61pub(crate) fn needs_probe(divergence: &Divergence) -> bool {
62 if !divergence.behind {
63 return false;
64 }
65 if divergence.truncated {
66 return true;
67 }
68 let ours: HashSet<&str> = divergence.ours.iter().map(String::as_str).collect();
69 divergence.theirs.iter().any(|path| ours.contains(path.as_str()))
70}
71
72/// Paths as a sandbox reported them, tidied: trimmed, without duplicates,
73/// and no more than are useful.
74pub(crate) fn tidy(paths: Vec<String>) -> Vec<String> {
75 let mut tidy: Vec<String> = Vec::new();
76 for path in paths {
77 let path = path.trim().to_owned();
78 if !path.is_empty() && !tidy.contains(&path) {
79 tidy.push(path);
80 }
81 if tidy.len() == MAX_CONFLICTS {
82 break;
83 }
84 }
85 tidy
86}
87
88impl Work {
89 async fn merge_row(&self, pull_id: &str) -> Result<Option<MergeRow>> {
90 self.db
91 .prepare(
92 "SELECT mergeable, conflicts, mergeable_key, mergeable_until
93 FROM pulls WHERE id = ?",
94 )
95 .bind(&[pull_id.into()])?
96 .first::<MergeRow>(None)
97 .await
98 }
99
100 /// Works out whether a pull request merges cleanly as it and its target
101 /// are now, unless that pair of commits was already worked out or is
102 /// being. Settles it at once where it can, and asks for a probe where
103 /// it cannot.
104 pub(crate) async fn assess_mergeability(&self, pull: &Pull) -> Result<()> {
105 if !pull.status.is_active() {
106 return Ok(());
107 }
108 let divergence: Option<Divergence> = g1t_kit::call(
109 &self.repos,
110 "divergence",
111 &BehindArgs {
112 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| pull.repo_id.clone()),
113 branch: pull.branch.clone(),
114 },
115 )
116 .await?;
117 // Nothing pushed yet.
118 let Some(divergence) = divergence else {
119 return Ok(());
120 };
121 let key = key(&divergence);
122 let now = now_ms();
123 let behind = u32::from(divergence.behind);
124 let row = self.merge_row(&pull.id).await?;
125 if let Some(row) = &row
126 && row.mergeable_key.as_deref() == Some(key.as_str())
127 {
128 // Whether it is behind goes with the pair of commits; a row from
129 // before it was kept gets it now.
130 self.db
131 .prepare("UPDATE pulls SET behind = ?1 WHERE id = ?2 AND behind IS NOT ?1")
132 .bind(&[behind.into(), pull.id.as_str().into()])?
133 .run()
134 .await?;
135 let settled = match row.mergeable.as_deref() {
136 Some("clean" | "conflicting" | "unknown") => true,
137 // A probe that is still being waited for.
138 Some("checking") => row
139 .mergeable_until
140 .as_deref()
141 .is_some_and(|until| until > rfc3339(now).as_str()),
142 _ => false,
143 };
144 if settled {
145 return Ok(());
146 }
147 }
148 if !needs_probe(&divergence) {
149 self.db
150 .prepare(
151 "UPDATE pulls
152 SET mergeable = 'clean', conflicts = '[]', mergeable_key = ?, behind = ?,
153 mergeable_token_hash = NULL, mergeable_until = NULL
154 WHERE id = ?",
155 )
156 .bind(&[key.as_str().into(), behind.into(), pull.id.as_str().into()])?
157 .run()
158 .await?;
159 // It may have been conflicting before this push.
160 if row.is_some_and(|row| row.mergeable.as_deref() == Some("conflicting")) {
161 self.announce_mergeability(pull).await?;
162 }
163 return Ok(());
164 }
165 self.db
166 .prepare(
167 "UPDATE pulls
168 SET mergeable = 'checking', conflicts = '[]', mergeable_key = ?, behind = ?,
169 mergeable_token_hash = NULL, mergeable_until = ?
170 WHERE id = ?",
171 )
172 .bind(&[
173 key.as_str().into(),
174 behind.into(),
175 rfc3339(now + PROBE_MINUTES * 60 * 1000).into(),
176 pull.id.as_str().into(),
177 ])?
178 .run()
179 .await?;
180 self.ask_for_probe(pull, &divergence.head).await
181 }
182
183 async fn ask_for_probe(&self, pull: &Pull, head: &str) -> Result<()> {
184 self.publish_as(
185 "pull.mergecheck",
186 &pull.repo_id,
187 None,
188 g1t_contracts::events::PullEvent {
189 commit: Some(head.to_owned()),
190 ..Self::pull_event(pull)
191 },
192 )
193 .await
194 }
195
196 /// Says that a pull request's mergeability settled, so that the runner
197 /// looks again at one g1t is seeing through: a conflict is the agent's
198 /// to resolve.
199 async fn announce_mergeability(&self, pull: &Pull) -> Result<()> {
200 self.publish_as("pull.mergeability", &pull.repo_id, None, Self::pull_event(pull))
201 .await
202 }
203
204 /// Works out mergeability again after a push: for the pull requests
205 /// whose heads it moved, and, when it moved a repository's default
206 /// branch, for every open pull request into it. A failure is logged and
207 /// not passed on, so that it never holds up the rest of the push.
208 pub(crate) async fn after_push(&self, repo_id: &str, default_branch: bool, moved: &[String]) {
209 let mut ids: Vec<String> = moved.to_vec();
210 if default_branch {
211 match self.targeting(repo_id).await {
212 Ok(targeting) => {
213 for id in targeting {
214 if !ids.contains(&id) {
215 ids.push(id);
216 }
217 }
218 }
219 Err(error) => worker::console_warn!("mergeability: {error}"),
220 }
221 }
222 for id in ids {
223 let assessed = async {
224 if let Some(pull) = self.pull_by_id(&id).await? {
225 self.assess_mergeability(&pull).await?;
226 }
227 Ok::<_, worker::Error>(())
228 };
229 if let Err(error) = assessed.await {
230 worker::console_warn!("mergeability of {id}: {error}");
231 }
232 }
233 }
234
235 /// The open pull requests into a repository, most recently active first.
236 async fn targeting(&self, repo_id: &str) -> Result<Vec<String>> {
237 Ok(self
238 .db
239 .prepare(
240 "SELECT id AS value FROM pulls
241 WHERE repo_id = ? AND status IN ('draft', 'open')
242 ORDER BY updated_at DESC LIMIT ?",
243 )
244 .bind(&[repo_id.into(), MAX_TARGETING.into()])?
245 .all()
246 .await?
247 .results::<ValueRow>()?
248 .into_iter()
249 .map(|row| row.value)
250 .collect())
251 }
252
253 /// Where a pull request's mergeability stands, for showing it. One that
254 /// was never worked out, or whose probe was given up on, is worked out
255 /// now. A conflict found for an older head is not shown as one.
256 pub(crate) async fn mergeability(&self, pull: &Pull) -> Result<(Mergeable, Vec<String>)> {
257 if !pull.status.is_active() || pull.head_commit.is_none() {
258 return Ok((Mergeable::Unknown, Vec::new()));
259 }
260 // As read for this request, if it was; read again after working it out.
261 let row = match self.prefetched_pull(&pull.id) {
262 Some(found) => found.first::<MergeRow>(crate::prefetch::Slot::Pull)?,
263 None => self.merge_row(&pull.id).await?,
264 };
265 let now = rfc3339(now_ms());
266 let stale = row.as_ref().is_none_or(|row| {
267 row.mergeable_key.is_none()
268 || (row.mergeable.as_deref() == Some("checking")
269 && row.mergeable_until.as_deref().is_none_or(|until| until < now.as_str()))
270 });
271 let row = if stale {
272 if let Err(error) = self.assess_mergeability(pull).await {
273 worker::console_warn!("mergeability of {}: {error}", pull.id);
274 }
275 // What was read for this request is out of date now.
276 self.keep_prefetched(None);
277 self.merge_row(&pull.id).await?
278 } else {
279 row
280 };
281 let Some(row) = row else {
282 return Ok((Mergeable::Unknown, Vec::new()));
283 };
284 let current = row
285 .mergeable_key
286 .as_deref()
287 .and_then(|key| key.split_once(".."))
288 .is_some_and(|(head, _)| Some(head) == pull.head_commit.as_deref());
289 let state = Mergeable::parse(row.mergeable.as_deref());
290 if state == Mergeable::Conflicting && !current {
291 return Ok((Mergeable::Checking, Vec::new()));
292 }
293 let conflicts = if state == Mergeable::Conflicting {
294 row.conflicts
295 .as_deref()
296 .and_then(|conflicts| serde_json::from_str(conflicts).ok())
297 .unwrap_or_default()
298 } else {
299 Vec::new()
300 };
301 Ok((state, conflicts))
302 }
303
304 /// The files that conflict, if the pull request as it is now is known
305 /// to conflict with its target as it is now.
306 pub(crate) async fn conflicting_files(&self, pull: &Pull) -> Result<Option<Vec<String>>> {
307 let row = match self.prefetched_pull(&pull.id) {
308 Some(found) => found.first::<MergeRow>(crate::prefetch::Slot::Pull)?,
309 None => self.merge_row(&pull.id).await?,
310 };
311 let Some(row) = row else {
312 return Ok(None);
313 };
314 if row.mergeable.as_deref() != Some("conflicting") {
315 return Ok(None);
316 }
317 let current = row
318 .mergeable_key
319 .as_deref()
320 .and_then(|key| key.split_once(".."))
321 .is_some_and(|(head, _)| Some(head) == pull.head_commit.as_deref());
322 if !current {
323 return Ok(None);
324 }
325 Ok(Some(
326 row.conflicts
327 .as_deref()
328 .and_then(|conflicts| serde_json::from_str(conflicts).ok())
329 .unwrap_or_default(),
330 ))
331 }
332
333 /// Claims the probe a `pull.mergecheck` event asked for, and returns
334 /// what a sandbox needs to carry it out.
335 pub(crate) async fn start_mergecheck(&self, a: StartMergecheckArgs) -> Result<Outcome<MergecheckJob>> {
336 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
337 return Ok(Outcome::fail(FailureCode::NotFound, "Pull request not found."));
338 };
339 if !pull.status.is_active() {
340 return Ok(Outcome::fail(
341 FailureCode::Conflict,
342 format!("This pull request is already {}.", pull.status.as_str()),
343 ));
344 }
345 let Some(row) = self.merge_row(&pull.id).await? else {
346 return Ok(Outcome::fail(FailureCode::NotFound, "Pull request not found."));
347 };
348 let Some((head, base)) = row
349 .mergeable_key
350 .as_deref()
351 .filter(|_| row.mergeable.as_deref() == Some("checking"))
352 .and_then(|key| key.split_once(".."))
353 .map(|(head, base)| (head.to_owned(), base.to_owned()))
354 else {
355 return Ok(Outcome::fail(
356 FailureCode::Conflict,
357 "Nothing is waiting to be checked for this pull request.",
358 ));
359 };
360 let now = now_ms();
361 let timestamp = rfc3339(now);
362 let running = self
363 .db
364 .prepare(
365 "SELECT count(*) AS n FROM pulls
366 WHERE repo_id = ? AND id != ? AND mergeable = 'checking'
367 AND mergeable_token_hash IS NOT NULL AND mergeable_until > ?",
368 )
369 .bind(&[pull.repo_id.as_str().into(), pull.id.as_str().into(), timestamp.as_str().into()])?
370 .first::<NumberRow>(None)
371 .await?
372 .map_or(0, |row| row.n);
373 if running >= MAX_PROBES {
374 // It waits, and starts when one of those reports.
375 return Ok(Outcome::fail(
376 FailureCode::Conflict,
377 "This repository has as many merge checks running as it may.",
378 ));
379 }
380 let token = new_token();
381 let claimed = self
382 .db
383 .prepare(
384 "UPDATE pulls SET mergeable_token_hash = ?, mergeable_until = ?
385 WHERE id = ? AND mergeable = 'checking' AND mergeable_key = ?
386 AND (mergeable_token_hash IS NULL OR mergeable_until < ?)
387 RETURNING id AS value",
388 )
389 .bind(&[
390 hash(&token).into(),
391 rfc3339(now + PROBE_MINUTES * 60 * 1000).into(),
392 pull.id.as_str().into(),
393 format!("{head}..{base}").into(),
394 timestamp.as_str().into(),
395 ])?
396 .first::<ValueRow>(None)
397 .await?;
398 if claimed.is_none() {
399 return Ok(Outcome::fail(
400 FailureCode::Conflict,
401 "This pull request is already being checked.",
402 ));
403 }
404 // Its author can read both the repository and the change.
405 let repo: Outcome<Repo> = g1t_kit::call(
406 &self.repos,
407 "get_by_id",
408 &GetByIdArgs {
409 id: pull.repo_id.clone(),
410 viewer: self.author_viewer(&pull).await?,
411 },
412 )
413 .await?;
414 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
415 return Ok(Outcome::fail(FailureCode::NotFound, "Pull request not found."));
416 };
417 let path = RepoPath {
418 namespace: repo.namespace.clone(),
419 name: repo.name.clone(),
420 };
421 Ok(Outcome::Ok(MergecheckJob {
422 pull_id: pull.id.clone(),
423 token,
424 source: pull.fork.clone().unwrap_or_else(|| path.clone()),
425 branch: pull.branch.clone().unwrap_or_else(|| repo.default_branch.clone()),
426 repo: path,
427 number: pull.number,
428 default_branch: repo.default_branch,
429 base,
430 head,
431 author: pull.author,
432 }))
433 }
434
435 /// Records what a probe found. Its token is the only credential, and a
436 /// probe for a pair of commits that has since moved on is refused.
437 pub(crate) async fn report_mergecheck(&self, a: ReportMergecheckArgs) -> Result<Outcome<Mergeable>> {
438 let row = self
439 .db
440 .prepare("SELECT id, repo_id, mergeable_key, mergeable_token_hash FROM pulls WHERE id = ?")
441 .bind(&[a.pull_id.as_str().into()])?
442 .first::<ProbeRow>(None)
443 .await?;
444 let Some(row) = row.filter(|row| row.mergeable_token_hash.as_deref() == Some(hash(&a.token).as_str())) else {
445 return Ok(Outcome::fail(FailureCode::NotFound, "Merge check not found."));
446 };
447 let conflicts = tidy(a.conflicts);
448 let state = if a.error.is_some() {
449 Mergeable::Unknown
450 } else if conflicts.is_empty() {
451 Mergeable::Clean
452 } else {
453 Mergeable::Conflicting
454 };
455 if let Some(error) = &a.error {
456 worker::console_warn!("merge check of {} could not run: {error}", row.id);
457 }
458 self.db
459 .prepare(
460 "UPDATE pulls
461 SET mergeable = ?, conflicts = ?, mergeable_token_hash = NULL, mergeable_until = NULL
462 WHERE id = ? AND mergeable_key IS ? AND mergeable_token_hash = ?",
463 )
464 .bind(&[
465 state.as_str().into(),
466 serde_json::to_string(&conflicts)?.into(),
467 row.id.as_str().into(),
468 crate::optional(&row.mergeable_key),
469 hash(&a.token).into(),
470 ])?
471 .run()
472 .await?;
473 if let Some(pull) = self.pull_by_id(&row.id).await? {
474 self.announce_mergeability(&pull).await?;
475 }
476 // The next one waiting its turn in this repository.
477 let waiting = self
478 .db
479 .prepare(
480 "SELECT id AS value FROM pulls
481 WHERE repo_id = ? AND mergeable = 'checking' AND mergeable_token_hash IS NULL
482 AND status IN ('draft', 'open')
483 ORDER BY updated_at DESC LIMIT 1",
484 )
485 .bind(&[row.repo_id.as_str().into()])?
486 .first::<ValueRow>(None)
487 .await?;
488 if let Some(waiting) = waiting
489 && let Some(pull) = self.pull_by_id(&waiting.value).await?
490 {
491 let head = pull.head_commit.clone().unwrap_or_default();
492 self.ask_for_probe(&pull, &head).await?;
493 }
494 Ok(Outcome::Ok(state))
495 }
496}
497
498#[cfg(test)]
499mod tests {
500 use super::*;
501
502 fn divergence(behind: bool, ours: &[&str], theirs: &[&str]) -> Divergence {
503 Divergence {
504 head: "h".into(),
505 base: "b".into(),
506 merge_base: Some("m".into()),
507 behind,
508 ours: ours.iter().map(|path| (*path).to_owned()).collect(),
509 theirs: theirs.iter().map(|path| (*path).to_owned()).collect(),
510 truncated: false,
511 }
512 }
513
514 #[test]
515 fn a_branch_that_has_not_moved_cannot_conflict() {
516 assert!(!needs_probe(&divergence(false, &["a.rs"], &[])));
517 }
518
519 #[test]
520 fn changes_to_different_files_cannot_conflict() {
521 assert!(!needs_probe(&divergence(true, &["a.rs", "b.rs"], &["c.rs"])));
522 }
523
524 #[test]
525 fn changes_to_the_same_file_are_probed() {
526 assert!(needs_probe(&divergence(true, &["a.rs", "b.rs"], &["b.rs"])));
527 }
528
529 #[test]
530 fn a_comparison_cut_short_is_probed() {
531 let cut = Divergence {
532 truncated: true,
533 ..divergence(true, &["a.rs"], &["c.rs"])
534 };
535 assert!(needs_probe(&cut));
536 }
537
538 #[test]
539 fn reported_paths_are_tidied() {
540 let paths = vec![" src/a.rs ".to_owned(), "src/a.rs".to_owned(), String::new(), "b.rs".to_owned()];
541 assert_eq!(tidy(paths), ["src/a.rs", "b.rs"]);
542 let many: Vec<String> = (0..150).map(|n| format!("f{n}")).collect();
543 assert_eq!(tidy(many).len(), MAX_CONFLICTS);
544 }
545}