Skip to content
1,785 linesCodeBlameRaw
1//! Mirroring: a repository's links to copies of it on other hosts (see
2//! `g1t_contracts::mirrors` for the model).
3//!
4//! Every linked repository has exactly one leader. A **mirror** follows a
5//! remote that leads: it stands by as an exact, read-only copy, and nothing
6//! runs on it. Someone can turn on **CI failover** (the remote keeps the
7//! code, g1t runs its workflows) or **take over** (g1t leads for a while).
8//! A takeover is **handed back** ref by ref: what only g1t changed is
9//! pushed, a protected branch goes as a pull request, and a branch both
10//! sides changed waits for a person's decision. A mirror can also be
11//! **moved to g1t** for good, after which g1t no longer tracks the remote.
12//! A repository g1t leads can be **mirrored to** any number of followers.
13//!
14//! This service keeps the links (`remotes`), each host's health
15//! (`remote_hosts`) and a takeover's starting point (`remote_refs`), and
16//! decides. The repos service keeps each repository's `RepoMirror`, set
17//! here with `set_mirror`, so pushes and merges are refused or allowed
18//! without asking. Git itself moves through the repos service (`mirror`,
19//! `mirror_refs`, `mirror_apply`).
20//!
21//! Hosts are adapters: [`RemoteProvider`] names them, and the few places a
22//! host's own API is used (credentials, whether a branch is protected,
23//! opening a pull request) match on it. GitHub goes through g1t's GitHub
24//! App; another g1t or any git host takes a username and token. A host
25//! that sends no webhook is polled.
26//!
27//! Nothing happens on its own unless someone asked for it in the link's
28//! settings: by default an unreachable remote is only shown on the
29//! repository, and a takeover starts when someone starts it.
30
31use std::collections::{BTreeMap, BTreeSet, HashMap};
32
33use g1t_contracts::access::{self, Capability};
34use g1t_contracts::events::{NewEvent, Publish};
35use g1t_contracts::identity::{ListMembersArgs, Member};
36use g1t_contracts::mirrors::*;
37use g1t_contracts::repos::{
38 GetByIdArgs, MirrorApplied, MirrorApplyArgs, MirrorArgs, MirrorDirection, MirrorRefs, MirrorRefsArgs, Mirrored, RefMove,
39 Repo, SetMirrorArgs,
40};
41use g1t_contracts::time::rfc3339;
42use g1t_contracts::{FailureCode, Outcome, Role, User, new_id};
43use g1t_kit::{args, now_ms, reply};
44use g1t_secrets::Sealer;
45use serde::Deserialize;
46use serde_json::Value;
47use worker::wasm_bindgen::JsValue;
48use worker::{Context, D1Database, Env, Fetcher, Method, Response, Result};
49
50use crate::github::GithubApp;
51use crate::http;
52
53const SOURCE: &str = "integrations";
54/// A host is unreachable after this many failed checks in a row, spanning
55/// at least [`DOWN_AFTER_MS`].
56const DOWN_CHECKS: u32 = 3;
57const DOWN_AFTER_MS: u64 = 2 * 60 * 1000;
58/// And reachable again after this many good ones, spanning [`UP_AFTER_MS`].
59const UP_CHECKS: u32 = 3;
60const UP_AFTER_MS: u64 = 5 * 60 * 1000;
61/// How often a remote with no webhook is asked for news.
62const POLL_MS: u64 = 5 * 60 * 1000;
63/// Where g1t's commits go when the remote will not take them directly.
64const HANDBACK_PREFIX: &str = "refs/heads/g1t/handback/";
65/// The most branches asked about protection in one plan.
66const MAX_PROTECTION_CHECKS: usize = 20;
67
68fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
69 Outcome::fail(code, message)
70}
71
72fn null_or(value: Option<&str>) -> JsValue {
73 value.map_or(JsValue::NULL, Into::into)
74}
75
76// --- Pure parts, tested below ----------------------------------------------
77
78/// The host of an https address, lowercased: `github.com`.
79pub fn host_of(url: &str) -> String {
80 url.trim()
81 .trim_start_matches("https://")
82 .trim_start_matches("http://")
83 .split(['/', '?', '#'])
84 .next()
85 .unwrap_or_default()
86 .rsplit('@')
87 .next()
88 .unwrap_or_default()
89 .to_lowercase()
90}
91
92/// A remote as people name it: `github.com/acme/web`.
93pub fn display_name(url: &str) -> String {
94 let rest = url
95 .trim()
96 .trim_start_matches("https://")
97 .trim_start_matches("http://")
98 .trim_end_matches('/');
99 let rest = rest.strip_suffix(".git").unwrap_or(rest);
100 let rest = rest.rsplit_once('@').map_or(rest, |(_, after)| after);
101 let (host, path) = rest.split_once('/').unwrap_or((rest, ""));
102 if path.is_empty() { host.to_lowercase() } else { format!("{}/{path}", host.to_lowercase()) }
103}
104
105/// A clone address: https, with `.git` added when the host leaves it off,
106/// and no credentials in it.
107pub fn clean_clone_url(url: &str) -> Option<String> {
108 let url = url.trim().trim_end_matches('/');
109 let rest = url.strip_prefix("https://")?;
110 if rest.contains('@') || rest.contains(' ') || !rest.contains('/') || rest.starts_with('/') {
111 return None;
112 }
113 Some(if url.ends_with(".git") { url.to_owned() } else { format!("{url}.git") })
114}
115
116/// The web address of a clone address.
117pub fn web_url(clone_url: &str) -> String {
118 clone_url.strip_suffix(".git").unwrap_or(clone_url).to_owned()
119}
120
121/// How a host's health moves with one more check. `now` in milliseconds.
122#[derive(Clone, Debug, Default, PartialEq, Eq)]
123pub struct HostHealth {
124 pub failures: u32,
125 pub successes: u32,
126 /// When the current run of failures or successes began.
127 pub streak_ms: u64,
128 pub unreachable_since: Option<u64>,
129}
130
131#[derive(Clone, Copy, Debug, PartialEq, Eq)]
132pub enum Change {
133 None,
134 WentDown,
135 CameBack,
136}
137
138impl HostHealth {
139 pub fn reachable(&self) -> bool {
140 self.unreachable_since.is_none()
141 }
142
143 pub fn checked(&self, answered: bool, now: u64) -> (HostHealth, Change) {
144 let mut next = self.clone();
145 if answered {
146 if next.successes == 0 {
147 next.streak_ms = now;
148 }
149 next.successes += 1;
150 next.failures = 0;
151 if next.unreachable_since.is_some()
152 && next.successes >= UP_CHECKS
153 && now.saturating_sub(next.streak_ms) >= UP_AFTER_MS
154 {
155 next.unreachable_since = None;
156 return (next, Change::CameBack);
157 }
158 } else {
159 if next.failures == 0 {
160 next.streak_ms = now;
161 }
162 next.failures += 1;
163 next.successes = 0;
164 if next.unreachable_since.is_none()
165 && next.failures >= DOWN_CHECKS
166 && now.saturating_sub(next.streak_ms) >= DOWN_AFTER_MS
167 {
168 next.unreachable_since = Some(now);
169 return (next, Change::WentDown);
170 }
171 }
172 (next, Change::None)
173 }
174}
175
176/// The branch name in a ref: `main` for `refs/heads/main`.
177fn branch_of(git_ref: &str) -> Option<&str> {
178 git_ref.strip_prefix("refs/heads/")
179}
180
181/// The moves that carry out a hand-back plan, and the refs that go as pull
182/// requests. `theirs_all` is every ref the remote has, so a hand-back
183/// branch made before is moved from where it is.
184pub fn hand_back_moves(plan: &HandbackPlan, theirs_all: &BTreeMap<String, String>) -> (Vec<RefMove>, Vec<String>) {
185 let mut moves = Vec::new();
186 let mut pull_requests = Vec::new();
187 for item in &plan.refs {
188 let push = |moves: &mut Vec<RefMove>| {
189 moves.push(RefMove {
190 direction: MirrorDirection::Push,
191 git_ref: item.git_ref.clone(),
192 to: None,
193 old: item.theirs.clone(),
194 new: item.ours.clone(),
195 })
196 };
197 let fetch = |moves: &mut Vec<RefMove>| {
198 moves.push(RefMove {
199 direction: MirrorDirection::Pull,
200 git_ref: item.git_ref.clone(),
201 to: None,
202 old: item.ours.clone(),
203 new: item.theirs.clone(),
204 })
205 };
206 let decision = match item.action {
207 RefAction::Diverged => item.decision,
208 RefAction::PullRequest => Some(RefDecision::PullRequest),
209 _ => None,
210 };
211 match (item.action, decision) {
212 (RefAction::Same, _) => {}
213 (RefAction::Push, _) | (RefAction::Diverged, Some(RefDecision::KeepOurs)) => push(&mut moves),
214 (RefAction::Fetch, _) | (RefAction::Diverged, Some(RefDecision::KeepTheirs)) => fetch(&mut moves),
215 (_, Some(RefDecision::PullRequest)) => {
216 let Some(branch) = branch_of(&item.git_ref) else {
217 // Only branches can be pulled; a tag goes as it is.
218 push(&mut moves);
219 continue;
220 };
221 if let Some(ours) = &item.ours {
222 let target = format!("{HANDBACK_PREFIX}{branch}");
223 moves.push(RefMove {
224 direction: MirrorDirection::Push,
225 git_ref: item.git_ref.clone(),
226 to: Some(target.clone()),
227 old: theirs_all.get(&target).cloned(),
228 new: Some(ours.clone()),
229 });
230 pull_requests.push(item.git_ref.clone());
231 }
232 // g1t follows the remote's branch; its commits are on the
233 // hand-back branch and kept under refs/g1t/replaced/.
234 fetch(&mut moves);
235 }
236 (RefAction::Diverged, None) | (RefAction::PullRequest, _) => {}
237 }
238 }
239 (moves, pull_requests)
240}
241
242/// Builds the plan for a hand-back from both sides' refs and what each ref
243/// was when the takeover began. `protected` names branches the remote
244/// protects.
245pub fn plan_refs(
246 bases: &BTreeMap<String, (Option<String>, Option<RefDecision>)>,
247 ours: &BTreeMap<String, String>,
248 theirs: &BTreeMap<String, String>,
249 protected: &BTreeSet<String>,
250) -> Vec<RefPlan> {
251 let names: BTreeSet<&String> = bases.keys().chain(ours.keys()).chain(theirs.keys()).collect();
252 names
253 .into_iter()
254 .filter(|name| !name.starts_with(HANDBACK_PREFIX))
255 .filter(|name| name.starts_with("refs/heads/") || name.starts_with("refs/tags/"))
256 .filter_map(|name| {
257 let (base, decision) = bases.get(name).cloned().unwrap_or((None, None));
258 let ours = ours.get(name).cloned();
259 let theirs = theirs.get(name).cloned();
260 // A ref g1t never had and the remote made since is copied in; a
261 // ref neither had at the start nor has now is nothing.
262 let action = ref_action(base.as_deref(), ours.as_deref(), theirs.as_deref(), protected.contains(name));
263 (action != RefAction::Same).then(|| RefPlan {
264 git_ref: name.clone(),
265 base,
266 ours,
267 theirs,
268 action,
269 decision: if action == RefAction::Diverged { decision } else { None },
270 })
271 })
272 .collect()
273}
274
275// --- Rows -------------------------------------------------------------------
276
277#[derive(Clone, Debug, Deserialize)]
278pub(crate) struct RemoteRow {
279 pub id: String,
280 pub repo_id: String,
281 pub workspace: String,
282 pub repo: String,
283 pub provider: String,
284 pub role: String,
285 pub name: String,
286 pub url: String,
287 pub clone_url: String,
288 pub connection_id: Option<String>,
289 pub username: Option<String>,
290 pub credential: Option<String>,
291 pub state: String,
292 pub state_since: String,
293 pub state_by: Option<String>,
294 pub settings: String,
295 pub synced_at: Option<String>,
296 pub last_error: Option<String>,
297 pub created_at: String,
298}
299
300impl RemoteRow {
301 fn provider(&self) -> RemoteProvider {
302 RemoteProvider::parse(&self.provider).unwrap_or(RemoteProvider::Git)
303 }
304
305 fn leads(&self) -> bool {
306 self.role == RemoteRole::Leader.as_str()
307 }
308
309 fn state(&self) -> RemoteState {
310 RemoteState::parse(&self.state)
311 }
312
313 fn settings(&self) -> MirrorSettings {
314 serde_json::from_str(&self.settings).unwrap_or_default()
315 }
316
317 fn host(&self) -> String {
318 host_of(&self.clone_url)
319 }
320
321 /// `owner/name` on GitHub.
322 fn full_name(&self) -> &str {
323 self.url.trim_start_matches("https://github.com/")
324 }
325
326 /// What the repos service keeps for a mirror of this remote.
327 fn repo_mirror(&self) -> Option<RepoMirror> {
328 if !self.leads() {
329 return None;
330 }
331 let settings = self.settings();
332 Some(RepoMirror {
333 state: self.state().mirror().unwrap_or_default(),
334 remote: self.name.clone(),
335 url: self.url.clone(),
336 since: self.state_since.clone(),
337 warm: settings.keep_ci_warm,
338 github_workflows: settings.github_workflows,
339 hold_deploys: settings.hold_deploys,
340 })
341 }
342
343 fn contract(&self, health: Option<&HostRow>) -> Remote {
344 Remote {
345 id: self.id.clone(),
346 repo_id: self.repo_id.clone(),
347 repo: self.repo.clone(),
348 provider: self.provider(),
349 role: if self.leads() { RemoteRole::Leader } else { RemoteRole::Follower },
350 name: self.name.clone(),
351 url: self.url.clone(),
352 state: self.state(),
353 state_since: self.state_since.clone(),
354 state_by: self.state_by.clone(),
355 reachable: health.is_none_or(|health| health.unreachable_since.is_none()),
356 unreachable_since: health.and_then(|health| health.unreachable_since.clone()),
357 synced_at: self.synced_at.clone(),
358 last_error: self.last_error.clone(),
359 settings: self.settings(),
360 created_at: self.created_at.clone(),
361 }
362 }
363
364 fn link(&self) -> String {
365 format!("/{}/settings/mirroring", self.repo)
366 }
367}
368
369#[derive(Clone, Debug, Default, Deserialize)]
370pub(crate) struct HostRow {
371 pub host: String,
372 pub failures: u32,
373 pub successes: u32,
374 pub streak_ms: f64,
375 pub unreachable_since: Option<String>,
376}
377
378impl HostRow {
379 fn health(&self) -> HostHealth {
380 HostHealth {
381 failures: self.failures,
382 successes: self.successes,
383 streak_ms: self.streak_ms as u64,
384 unreachable_since: self.unreachable_since.as_deref().and_then(crate::github::parse_time),
385 }
386 }
387}
388
389/// Who is acting: a person, or g1t on its own (an automatic takeover or
390/// hand-back), which needs no role.
391enum By<'a> {
392 Person(&'a User),
393 G1t,
394}
395
396impl By<'_> {
397 fn name(&self) -> String {
398 match self {
399 By::Person(user) => user.username.clone(),
400 By::G1t => "g1t".to_owned(),
401 }
402 }
403}
404
405pub struct Mirrors {
406 db: D1Database,
407 sealer: Option<Sealer>,
408 repos: Fetcher,
409 identity: Fetcher,
410 events: Option<Fetcher>,
411 github: GithubApp,
412}
413
414impl Mirrors {
415 pub fn new(env: &Env) -> Result<Self> {
416 Ok(Mirrors {
417 db: env.d1("DB")?,
418 sealer: env.secret("INTEGRATIONS_KEY").ok().and_then(|key| Sealer::new(&key.to_string())),
419 repos: env.service("REPOS")?,
420 identity: env.service("IDENTITY")?,
421 events: env.service("EVENTS").ok(),
422 github: GithubApp::new(env)?,
423 })
424 }
425
426 // --- Reading --------------------------------------------------------------
427
428 async fn rows(&self, repo_id: &str) -> Result<Vec<RemoteRow>> {
429 self.db
430 .prepare("SELECT * FROM remotes WHERE repo_id = ? ORDER BY role, created_at")
431 .bind(&[repo_id.into()])?
432 .all()
433 .await?
434 .results::<RemoteRow>()
435 }
436
437 async fn row(&self, id: &str) -> Result<Option<RemoteRow>> {
438 self.db.prepare("SELECT * FROM remotes WHERE id = ?").bind(&[id.into()])?.first::<RemoteRow>(None).await
439 }
440
441 async fn leader(&self, repo_id: &str) -> Result<Option<RemoteRow>> {
442 self.db
443 .prepare("SELECT * FROM remotes WHERE repo_id = ? AND role = 'leader'")
444 .bind(&[repo_id.into()])?
445 .first::<RemoteRow>(None)
446 .await
447 }
448
449 async fn hosts(&self) -> Result<HashMap<String, HostRow>> {
450 Ok(self
451 .db
452 .prepare("SELECT * FROM remote_hosts")
453 .all()
454 .await?
455 .results::<HostRow>()?
456 .into_iter()
457 .map(|row| (row.host.clone(), row))
458 .collect())
459 }
460
461 async fn repo(&self, repo_id: &str, viewer: Option<&User>) -> Result<Option<Repo>> {
462 let viewer = viewer.cloned().or_else(|| Some(User::system("g1t")));
463 let found: Outcome<Repo> =
464 g1t_kit::call(&self.repos, "get_by_id", &GetByIdArgs { id: repo_id.to_owned(), viewer }).await?;
465 Ok(found.into_result().ok())
466 }
467
468 /// Why `actor` may not change `repo`'s links, if they may not.
469 fn refused<T>(actor: &User, repo: &Repo, capability: Capability) -> Option<Outcome<T>> {
470 let target = access::RepoRef { id: &repo.id, namespace: &repo.namespace, private: repo.is_private };
471 match access::check(Some(actor), target, capability) {
472 Ok(()) => None,
473 Err(access::Denied::NotFound) => Some(fail(FailureCode::NotFound, "Repository not found.")),
474 Err(access::Denied::Forbidden) => {
475 Some(fail(FailureCode::Forbidden, access::needs(capability, &format!("{}/{}", repo.namespace, repo.name))))
476 }
477 }
478 }
479
480 /// The repository and its leader, for someone changing them.
481 async fn leader_for(&self, actor: &User, repo_id: &str) -> Result<std::result::Result<(Repo, RemoteRow), Outcome<MirrorView>>> {
482 let Some(repo) = self.repo(repo_id, Some(actor)).await? else {
483 return Ok(Err(fail(FailureCode::NotFound, "Repository not found.")));
484 };
485 if let Some(refused) = Self::refused(actor, &repo, Capability::ManageIntegrations) {
486 return Ok(Err(refused));
487 }
488 let Some(row) = self.leader(repo_id).await? else {
489 return Ok(Err(fail(FailureCode::NotFound, format!("{}/{} is not a mirror.", repo.namespace, repo.name))));
490 };
491 Ok(Ok((repo, row)))
492 }
493
494 async fn view_of(&self, repo_id: &str, can_manage: bool, plan: Option<HandbackPlan>) -> Result<MirrorView> {
495 let hosts = self.hosts().await?;
496 Ok(MirrorView {
497 remotes: self.rows(repo_id).await?.iter().map(|row| row.contract(hosts.get(&row.host()))).collect(),
498 plan,
499 can_manage,
500 notes: Vec::new(),
501 })
502 }
503
504 pub async fn view(&self, a: MirrorViewArgs) -> Result<Outcome<MirrorView>> {
505 let Some(repo) = self.repo(&a.repo_id, a.viewer.as_ref()).await? else {
506 return Ok(fail(FailureCode::NotFound, "Repository not found."));
507 };
508 let can_manage = a
509 .viewer
510 .as_ref()
511 .is_some_and(|viewer| Self::refused::<()>(viewer, &repo, Capability::ManageIntegrations).is_none());
512 Ok(Outcome::Ok(self.view_of(&repo.id, can_manage, None).await?))
513 }
514
515 pub async fn briefs(&self, a: MirrorBriefsArgs) -> Result<Vec<RemoteBrief>> {
516 let ids: Vec<&String> = a.repo_ids.iter().take(200).collect();
517 if ids.is_empty() {
518 return Ok(Vec::new());
519 }
520 let marks = vec!["?"; ids.len()].join(", ");
521 let bind: Vec<JsValue> = ids.iter().map(|id| id.as_str().into()).collect();
522 let rows = self
523 .db
524 .prepare(format!("SELECT * FROM remotes WHERE repo_id IN ({marks}) ORDER BY role, created_at"))
525 .bind(&bind)?
526 .all()
527 .await?
528 .results::<RemoteRow>()?;
529 let hosts = self.hosts().await?;
530 Ok(rows
531 .iter()
532 .map(|row| RemoteBrief {
533 repo_id: row.repo_id.clone(),
534 role: if row.leads() { RemoteRole::Leader } else { RemoteRole::Follower },
535 name: row.name.clone(),
536 state: row.state(),
537 reachable: hosts.get(&row.host()).is_none_or(|host| host.unreachable_since.is_none()),
538 synced_at: row.synced_at.clone(),
539 })
540 .collect())
541 }
542
543 // --- Talking to hosts -----------------------------------------------------
544
545 /// The user and token a remote is opened with.
546 async fn credential(&self, row: &RemoteRow) -> Result<std::result::Result<(Option<String>, String), String>> {
547 match row.provider() {
548 RemoteProvider::Github => {
549 let Some(installation) = row.connection_id.as_deref().and_then(|id| id.parse::<u64>().ok()) else {
550 return Ok(Err("This link has no GitHub installation.".to_owned()));
551 };
552 Ok(self.github.installation_token(installation).await?.map(|token| (None, token)))
553 }
554 RemoteProvider::G1t | RemoteProvider::Git => {
555 let (Some(sealer), Some(sealed)) = (&self.sealer, row.credential.as_deref()) else {
556 return Ok(Err("This link has no token. Add one in its settings.".to_owned()));
557 };
558 match sealer.open(sealed, &format!("rmt:{}", row.id)) {
559 Some(token) => Ok(Ok((row.username.clone(), token))),
560 None => Ok(Err("This link's token could not be read. Add it again.".to_owned())),
561 }
562 }
563 }
564 }
565
566 /// Which of `branches` the remote protects. Hosts that cannot say
567 /// protect none: a refused push still goes as a pull request.
568 async fn protected(&self, row: &RemoteRow, token: &str, branches: &[String]) -> Result<BTreeSet<String>> {
569 let mut protected = BTreeSet::new();
570 if row.provider() != RemoteProvider::Github {
571 return Ok(protected);
572 }
573 for git_ref in branches.iter().take(MAX_PROTECTION_CHECKS) {
574 let Some(branch) = branch_of(git_ref) else { continue };
575 let Ok(answer) = self
576 .github
577 .api(Method::Get, &format!("/repos/{}/branches/{}", row.full_name(), branch), token, None)
578 .await
579 else {
580 continue;
581 };
582 if answer.ok() && answer.json()["protected"].as_bool() == Some(true) {
583 protected.insert(git_ref.clone());
584 }
585 }
586 Ok(protected)
587 }
588
589 /// Opens a pull request on the remote from a hand-back branch. Returns
590 /// a line for people: where it is, or how to open it by hand.
591 async fn open_pull_request(&self, row: &RemoteRow, token: &str, git_ref: &str, repo: &str) -> Result<String> {
592 let branch = branch_of(git_ref).unwrap_or(git_ref);
593 let head = format!("g1t/handback/{branch}");
594 let by_hand = format!("{branch}: g1t's commits are on {head} on {}. Open a pull request from it into {branch}.", row.name);
595 if row.provider() != RemoteProvider::Github {
596 return Ok(by_hand);
597 }
598 let body = serde_json::json!({
599 "title": format!("Changes made on g1t while {} was unreachable", row.name),
600 "head": head,
601 "base": branch,
602 "body": format!(
603 "g1t led [{repo}](https://g1t.sh/{repo}) while this repository could not be reached. These are the commits made there on `{branch}`.\n\nThe branch is protected here, so they come as a pull request rather than a push."
604 ),
605 "maintainer_can_modify": true,
606 });
607 let Ok(answer) = self.github.api(Method::Post, &format!("/repos/{}/pulls", row.full_name()), token, Some(body)).await
608 else {
609 return Ok(by_hand);
610 };
611 let said = answer.json();
612 Ok(if answer.ok() {
613 format!("{branch}: opened {}", said["html_url"].as_str().unwrap_or("a pull request"))
614 } else if answer.status == 422 && said.to_string().contains("already exists") {
615 format!("{branch}: the pull request from {head} was already open and now has the new commits.")
616 } else {
617 format!("{by_hand} ({})", answer.problem("GitHub"))
618 })
619 }
620
621 async fn repos_refs(&self, row: &RemoteRow, username: Option<String>, token: String) -> Result<Outcome<MirrorRefs>> {
622 g1t_kit::call(
623 &self.repos,
624 "mirror_refs",
625 &MirrorRefsArgs { repo_id: row.repo_id.clone(), url: row.clone_url.clone(), token, username },
626 )
627 .await
628 }
629
630 /// g1t's own refs, asking nothing of the remote: a takeover starts while
631 /// it is away, when not even a token can be had from it.
632 async fn our_refs(&self, row: &RemoteRow) -> Result<Outcome<MirrorRefs>> {
633 g1t_kit::call(
634 &self.repos,
635 "mirror_refs",
636 &MirrorRefsArgs { repo_id: row.repo_id.clone(), url: String::new(), token: String::new(), username: None },
637 )
638 .await
639 }
640
641 async fn apply(&self, row: &RemoteRow, username: Option<String>, token: String, moves: Vec<RefMove>) -> Result<Outcome<MirrorApplied>> {
642 g1t_kit::call(
643 &self.repos,
644 "mirror_apply",
645 &MirrorApplyArgs { repo_id: row.repo_id.clone(), url: row.clone_url.clone(), token, username, moves },
646 )
647 .await
648 }
649
650 /// Copies refs in the link's direction: a mirror catches up, a follower
651 /// is pushed to. `Err` is a reason, already recorded on the link.
652 async fn sync(&self, row: &RemoteRow) -> Result<std::result::Result<Mirrored, String>> {
653 let direction = if row.leads() { MirrorDirection::Pull } else { MirrorDirection::Push };
654 let (username, token) = match self.credential(row).await? {
655 Ok(credential) => credential,
656 Err(reason) => {
657 self.note(row, Some(&reason)).await?;
658 return Ok(Err(reason));
659 }
660 };
661 let done: Outcome<Mirrored> = g1t_kit::call(
662 &self.repos,
663 "mirror",
664 &MirrorArgs { repo_id: row.repo_id.clone(), url: row.clone_url.clone(), token, username, direction },
665 )
666 .await?;
667 match done {
668 Outcome::Ok(mirrored) => {
669 self.note(row, None).await?;
670 self.host_checked(&row.host(), true, None).await?;
671 Ok(Ok(mirrored))
672 }
673 Outcome::Fail(failure) => {
674 let unreachable = failure.code == FailureCode::Unavailable;
675 if unreachable {
676 self.host_checked(&row.host(), false, Some(&failure.message)).await?;
677 } else {
678 self.note(row, Some(&failure.message)).await?;
679 }
680 if !row.leads() && !unreachable {
681 self.set_state(row, RemoteState::Stuck, None).await?;
682 }
683 Ok(Err(failure.message))
684 }
685 }
686 }
687
688 async fn note(&self, row: &RemoteRow, problem: Option<&str>) -> Result<()> {
689 let statement = match problem {
690 Some(problem) => self
691 .db
692 .prepare("UPDATE remotes SET last_error = ?2 WHERE id = ?1")
693 .bind(&[row.id.as_str().into(), problem.into()])?,
694 None => self
695 .db
696 .prepare("UPDATE remotes SET last_error = NULL, synced_at = ?2 WHERE id = ?1")
697 .bind(&[row.id.as_str().into(), rfc3339(now_ms()).into()])?,
698 };
699 statement.run().await?;
700 if problem.is_none() && row.state() == RemoteState::Stuck {
701 self.set_state(row, RemoteState::Following, None).await?;
702 }
703 Ok(())
704 }
705
706 // --- State ------------------------------------------------------------------
707
708 /// Records a link's state, and tells the repos service what a mirror
709 /// is now. A link the repos service has not heard of yet is retried by
710 /// the cron (`recorded`).
711 async fn set_state(&self, row: &RemoteRow, state: RemoteState, by: Option<&str>) -> Result<RemoteRow> {
712 let now = rfc3339(now_ms());
713 self.db
714 .prepare("UPDATE remotes SET state = ?2, state_since = ?3, state_by = ?4, recorded = 0 WHERE id = ?1")
715 .bind(&[row.id.as_str().into(), state.as_str().into(), now.as_str().into(), null_or(by)])?
716 .run()
717 .await?;
718 let row = RemoteRow {
719 state: state.as_str().to_owned(),
720 state_since: now,
721 state_by: by.map(str::to_owned),
722 ..row.clone()
723 };
724 self.record(&row).await?;
725 Ok(row)
726 }
727
728 /// Tells the repos service a repository's mirror, as this link has it.
729 async fn record(&self, row: &RemoteRow) -> Result<()> {
730 if !row.leads() {
731 return Ok(());
732 }
733 let done: Outcome<Repo> = g1t_kit::call(
734 &self.repos,
735 "set_mirror",
736 &SetMirrorArgs { repo_id: row.repo_id.clone(), mirror: row.repo_mirror() },
737 )
738 .await?;
739 if done.into_result().is_ok() {
740 self.db.prepare("UPDATE remotes SET recorded = 1 WHERE id = ?").bind(&[row.id.as_str().into()])?.run().await?;
741 }
742 Ok(())
743 }
744
745 async fn publish(&self, kind: &'static str, row: &RemoteRow, by: Option<&str>, title: String, detail: Option<String>, notify: Vec<String>) {
746 let Some(events) = &self.events else { return };
747 let event = NewEvent {
748 kind,
749 source: SOURCE,
750 repo_id: Some(row.repo_id.clone()),
751 actor: None,
752 data: MirrorEvent {
753 repo_id: row.repo_id.clone(),
754 repo: row.repo.clone(),
755 remote_id: row.id.clone(),
756 remote: row.name.clone(),
757 state: (kind != "mirror.moved_in").then(|| row.state.clone()),
758 from: None,
759 by: by.map(str::to_owned),
760 title,
761 detail,
762 notify,
763 link: row.link(),
764 },
765 };
766 let sent: Result<Value> = g1t_kit::call(events, "publish", &Publish { events: vec![event] }).await;
767 if let Err(error) = sent {
768 worker::console_error!("mirrors: publishing {kind} for {} failed: {error}", row.repo);
769 }
770 }
771
772 /// The workspace's owners, when the link asks for the inbox.
773 async fn told(&self, row: &RemoteRow, except: Option<&str>) -> Vec<String> {
774 if row.settings().notify != Notify::Inbox {
775 return Vec::new();
776 }
777 let members: Result<Outcome<Vec<Member>>> = g1t_kit::call(
778 &self.identity,
779 "list_members",
780 &ListMembersArgs { slug: row.workspace.clone(), viewer: Some(User::system(&row.workspace)) },
781 )
782 .await;
783 match members {
784 Ok(Outcome::Ok(members)) => members
785 .into_iter()
786 .filter(|member| member.role == Role::Owner)
787 .map(|member| member.username)
788 .filter(|name| except.is_none_or(|except| !name.eq_ignore_ascii_case(except)))
789 .collect(),
790 _ => Vec::new(),
791 }
792 }
793
794 // --- Takeover, CI failover, hand-back, moving in ------------------------------
795
796 pub async fn take_over(&self, a: MirrorActArgs) -> Result<Outcome<MirrorView>> {
797 let (_, row) = match self.leader_for(&a.actor, &a.repo_id).await? {
798 Ok(found) => found,
799 Err(refused) => return Ok(refused),
800 };
801 self.start_takeover(row, By::Person(&a.actor)).await
802 }
803
804 async fn start_takeover(&self, row: RemoteRow, by: By<'_>) -> Result<Outcome<MirrorView>> {
805 match row.state() {
806 RemoteState::Takeover => return Ok(Outcome::Ok(self.view_of(&row.repo_id, true, None).await?)),
807 RemoteState::HandingBack => {
808 return Ok(fail(FailureCode::Conflict, "It is being handed back. Wait for that to finish, or for it to stop."));
809 }
810 _ => {}
811 }
812 // The starting point: what g1t holds now, which is what it last
813 // copied from the remote.
814 let refs = match self.our_refs(&row).await? {
815 Outcome::Ok(refs) => refs,
816 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
817 };
818 let mut statements = vec![self.db.prepare("DELETE FROM remote_refs WHERE remote_id = ?").bind(&[row.id.as_str().into()])?];
819 for (name, hash) in &refs.ours {
820 statements.push(
821 self.db
822 .prepare("INSERT INTO remote_refs (remote_id, ref, base) VALUES (?, ?, ?)")
823 .bind(&[row.id.as_str().into(), name.as_str().into(), hash.as_str().into()])?,
824 );
825 }
826 self.db.batch(statements).await?;
827 let name = by.name();
828 let row = self.set_state(&row, RemoteState::Takeover, Some(&name)).await?;
829 let notify = self.told(&row, Some(&name)).await;
830 let title = match by {
831 By::G1t => format!("g1t took over {} while {} is unreachable", row.repo, row.name),
832 By::Person(_) => format!("{name} took over {} from {}", row.repo, row.name),
833 };
834 self.publish("mirror.state_changed", &row, Some(&name), title, None, notify).await;
835 Ok(Outcome::Ok(self.view_of(&row.repo_id, true, None).await?))
836 }
837
838 pub async fn ci(&self, a: MirrorCiArgs) -> Result<Outcome<MirrorView>> {
839 let (_, row) = match self.leader_for(&a.actor, &a.repo_id).await? {
840 Ok(found) => found,
841 Err(refused) => return Ok(refused),
842 };
843 let next = match (row.state(), a.on) {
844 (RemoteState::Standby, true) => RemoteState::Ci,
845 (RemoteState::Ci, false) => RemoteState::Standby,
846 (RemoteState::Ci, true) | (RemoteState::Standby, false) => {
847 return Ok(Outcome::Ok(self.view_of(&row.repo_id, true, None).await?));
848 }
849 _ => return Ok(fail(FailureCode::Conflict, "g1t leads this repository now: its workflows already run here.")),
850 };
851 let row = self.set_state(&row, next, Some(&a.actor.username)).await?;
852 let title = if a.on {
853 format!("{} started running {}'s workflows on g1t", a.actor.username, row.name)
854 } else {
855 format!("{} ended CI failover for {}", a.actor.username, row.repo)
856 };
857 self.publish("mirror.state_changed", &row, Some(&a.actor.username), title, None, Vec::new()).await;
858 Ok(Outcome::Ok(self.view_of(&row.repo_id, true, None).await?))
859 }
860
861 /// What handing back would do now.
862 async fn plan(&self, row: &RemoteRow) -> Result<Outcome<HandbackPlan>> {
863 let (username, token) = match self.credential(row).await? {
864 Ok(credential) => credential,
865 Err(reason) => return Ok(fail(FailureCode::Conflict, reason)),
866 };
867 let refs = match self.repos_refs(row, username, token.clone()).await? {
868 Outcome::Ok(refs) => refs,
869 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
870 };
871 let bases = self.bases(row).await?;
872 let Some(theirs) = refs.theirs else {
873 // The remote is away: what g1t changed, as it would go.
874 let refs = plan_refs(&bases, &refs.ours, &bases_as_theirs(&bases), &BTreeSet::new());
875 return Ok(Outcome::Ok(HandbackPlan::new(refs, false)));
876 };
877 let candidates: Vec<String> = plan_refs(&bases, &refs.ours, &theirs, &BTreeSet::new())
878 .into_iter()
879 .filter(|plan| plan.action == RefAction::Push && plan.ours.is_some() && plan.theirs.is_some())
880 .map(|plan| plan.git_ref)
881 .collect();
882 let protected = self.protected(row, &token, &candidates).await?;
883 Ok(Outcome::Ok(HandbackPlan::new(plan_refs(&bases, &refs.ours, &theirs, &protected), true)))
884 }
885
886 async fn bases(&self, row: &RemoteRow) -> Result<BTreeMap<String, (Option<String>, Option<RefDecision>)>> {
887 #[derive(Deserialize)]
888 struct Base {
889 #[serde(rename = "ref")]
890 git_ref: String,
891 base: Option<String>,
892 decision: Option<String>,
893 }
894 Ok(self
895 .db
896 .prepare("SELECT ref, base, decision FROM remote_refs WHERE remote_id = ?")
897 .bind(&[row.id.as_str().into()])?
898 .all()
899 .await?
900 .results::<Base>()?
901 .into_iter()
902 .map(|base| {
903 let decision = base.decision.and_then(|text| serde_json::from_value(Value::String(text)).ok());
904 (base.git_ref, (base.base, decision))
905 })
906 .collect())
907 }
908
909 pub async fn hand_back_plan(&self, a: MirrorActArgs) -> Result<Outcome<MirrorView>> {
910 let (_, row) = match self.leader_for(&a.actor, &a.repo_id).await? {
911 Ok(found) => found,
912 Err(refused) => return Ok(refused),
913 };
914 if !matches!(row.state(), RemoteState::Takeover | RemoteState::HandingBack) {
915 return Ok(fail(FailureCode::Conflict, "g1t has not taken over, so there is nothing to hand back."));
916 }
917 let plan = match self.plan(&row).await? {
918 Outcome::Ok(plan) => plan,
919 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
920 };
921 Ok(Outcome::Ok(self.view_of(&row.repo_id, true, Some(plan)).await?))
922 }
923
924 pub async fn hand_back(&self, a: MirrorHandBackArgs) -> Result<Outcome<MirrorView>> {
925 let (_, row) = match self.leader_for(&a.actor, &a.repo_id).await? {
926 Ok(found) => found,
927 Err(refused) => return Ok(refused),
928 };
929 let mut statements = Vec::new();
930 for (git_ref, decision) in &a.decisions {
931 let decision = serde_json::to_value(decision).ok().and_then(|value| value.as_str().map(str::to_owned));
932 statements.push(
933 self.db
934 .prepare(
935 "INSERT INTO remote_refs (remote_id, ref, decision) VALUES (?1, ?2, ?3)
936 ON CONFLICT (remote_id, ref) DO UPDATE SET decision = excluded.decision",
937 )
938 .bind(&[row.id.as_str().into(), git_ref.as_str().into(), null_or(decision.as_deref())])?,
939 );
940 }
941 if !statements.is_empty() {
942 self.db.batch(statements).await?;
943 }
944 self.finish_takeover(row, By::Person(&a.actor)).await
945 }
946
947 /// Hands a takeover back: freezes the repository, moves each ref as the
948 /// plan says, and returns to standing by. If anything is refused, g1t
949 /// keeps the lead and says why, so nobody is stuck.
950 async fn finish_takeover(&self, row: RemoteRow, by: By<'_>) -> Result<Outcome<MirrorView>> {
951 if !matches!(row.state(), RemoteState::Takeover | RemoteState::HandingBack) {
952 return Ok(fail(FailureCode::Conflict, "g1t has not taken over, so there is nothing to hand back."));
953 }
954 let plan = match self.plan(&row).await? {
955 Outcome::Ok(plan) => plan,
956 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
957 };
958 if !plan.reachable {
959 return Ok(fail(FailureCode::Unavailable, format!("{} is not answering yet. Hand back once it is.", row.name)));
960 }
961 if !plan.ready {
962 return Ok(Outcome::Fail(undecided(&plan)));
963 }
964 let name = by.name();
965 // Nothing moves on g1t while it goes back.
966 let row = self.set_state(&row, RemoteState::HandingBack, Some(&name)).await?;
967 let plan = match self.plan(&row).await? {
968 Outcome::Ok(plan) if plan.ready => plan,
969 Outcome::Ok(plan) => {
970 let row = self.set_state(&row, RemoteState::Takeover, row.state_by.as_deref()).await?;
971 let _ = row;
972 return Ok(Outcome::Fail(undecided(&plan)));
973 }
974 Outcome::Fail(failure) => {
975 self.set_state(&row, RemoteState::Takeover, Some(&name)).await?;
976 return Ok(Outcome::Fail(failure));
977 }
978 };
979 let (username, token) = match self.credential(&row).await? {
980 Ok(credential) => credential,
981 Err(reason) => {
982 self.set_state(&row, RemoteState::Takeover, Some(&name)).await?;
983 return Ok(fail(FailureCode::Conflict, reason));
984 }
985 };
986 let theirs_all = match self.repos_refs(&row, username.clone(), token.clone()).await? {
987 Outcome::Ok(refs) => refs.theirs.unwrap_or_default(),
988 Outcome::Fail(_) => BTreeMap::new(),
989 };
990 let (moves, mut pull_requests) = hand_back_moves(&plan, &theirs_all);
991 let mut problems = Vec::new();
992 if !moves.is_empty() {
993 let applied = match self.apply(&row, username.clone(), token.clone(), moves).await? {
994 Outcome::Ok(applied) => applied,
995 Outcome::Fail(failure) => {
996 self.set_state(&row, RemoteState::Takeover, Some(&name)).await?;
997 self.note(&row, Some(&failure.message)).await?;
998 return Ok(Outcome::Fail(failure));
999 }
1000 };
1001 // A push the remote refused (a protected branch it did not say
1002 // was protected) goes as a pull request instead.
1003 let refused: Vec<RefPlan> = applied
1004 .moved
1005 .iter()
1006 .filter(|moved| moved.direction == MirrorDirection::Push && moved.problem.is_some())
1007 .filter_map(|moved| {
1008 plan.refs
1009 .iter()
1010 .find(|item| item.git_ref == moved.git_ref && branch_of(&item.git_ref).is_some() && item.ours.is_some())
1011 .cloned()
1012 })
1013 .collect();
1014 for moved in applied.moved.iter().filter(|moved| moved.problem.is_some()) {
1015 if !refused.iter().any(|item| item.git_ref == moved.git_ref) {
1016 problems.push(format!("{}: {}", moved.git_ref, moved.problem.as_deref().unwrap_or_default()));
1017 }
1018 }
1019 if !refused.is_empty() {
1020 let retry = HandbackPlan::new(
1021 refused
1022 .into_iter()
1023 .map(|item| RefPlan { action: RefAction::PullRequest, ..item })
1024 .collect(),
1025 true,
1026 );
1027 let (moves, more) = hand_back_moves(&retry, &theirs_all);
1028 match self.apply(&row, username.clone(), token.clone(), moves).await? {
1029 Outcome::Ok(applied) => {
1030 for moved in applied.moved.iter().filter(|moved| moved.problem.is_some()) {
1031 problems.push(format!("{}: {}", moved.git_ref, moved.problem.as_deref().unwrap_or_default()));
1032 }
1033 pull_requests.extend(more);
1034 }
1035 Outcome::Fail(failure) => problems.push(failure.message),
1036 }
1037 }
1038 }
1039 if !problems.is_empty() {
1040 let reason = format!("Some branches did not go back: {}", problems.join("; "));
1041 let row = self.set_state(&row, RemoteState::Takeover, Some(&name)).await?;
1042 self.note(&row, Some(&reason)).await?;
1043 return Ok(fail(FailureCode::Conflict, format!("{reason}. g1t still leads; try again or decide differently.")));
1044 }
1045 let mut notes = Vec::new();
1046 for git_ref in &pull_requests {
1047 notes.push(self.open_pull_request(&row, &token, git_ref, &row.repo).await?);
1048 }
1049 self.db.prepare("DELETE FROM remote_refs WHERE remote_id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1050 let row = self.set_state(&row, RemoteState::Standby, Some(&name)).await?;
1051 // Anything else the remote has, g1t now follows.
1052 if let Err(reason) = self.sync(&row).await? {
1053 notes.push(format!("Catching up after handing back: {reason}"));
1054 }
1055 let notify = self.told(&row, Some(&name)).await;
1056 let detail = (!notes.is_empty()).then(|| notes.join("\n"));
1057 self.publish(
1058 "mirror.state_changed",
1059 &row,
1060 Some(&name),
1061 format!("{} was handed back to {}", row.repo, row.name),
1062 detail,
1063 notify,
1064 )
1065 .await;
1066 let mut view = self.view_of(&row.repo_id, true, None).await?;
1067 view.notes = notes;
1068 Ok(Outcome::Ok(view))
1069 }
1070
1071 pub async fn move_in(&self, a: MirrorMoveInArgs) -> Result<Outcome<MirrorView>> {
1072 let (_, row) = match self.leader_for(&a.actor, &a.repo_id).await? {
1073 Ok(found) => found,
1074 Err(refused) => return Ok(refused),
1075 };
1076 if row.state() == RemoteState::HandingBack {
1077 return Ok(fail(FailureCode::Conflict, "It is being handed back. Move it to g1t once that is done."));
1078 }
1079 let done: Outcome<Repo> =
1080 g1t_kit::call(&self.repos, "set_mirror", &SetMirrorArgs { repo_id: row.repo_id.clone(), mirror: None }).await?;
1081 if let Outcome::Fail(failure) = done {
1082 return Ok(Outcome::Fail(failure));
1083 }
1084 self.db.prepare("DELETE FROM remote_refs WHERE remote_id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1085 let mut notes = Vec::new();
1086 if a.keep_remote_updated {
1087 let now = rfc3339(now_ms());
1088 self.db
1089 .prepare(
1090 "UPDATE remotes SET role = 'follower', state = 'following', state_since = ?2, state_by = ?3, recorded = 1
1091 WHERE id = ?1",
1092 )
1093 .bind(&[row.id.as_str().into(), now.as_str().into(), a.actor.username.as_str().into()])?
1094 .run()
1095 .await?;
1096 if let Some(follower) = self.row(&row.id).await?
1097 && let Err(reason) = self.sync(&follower).await?
1098 {
1099 notes.push(format!("{} could not be brought up to date yet: {reason}", row.name));
1100 }
1101 } else {
1102 self.db.prepare("DELETE FROM remotes WHERE id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1103 }
1104 if row.provider() == RemoteProvider::Github {
1105 let mode = if a.keep_remote_updated { "push" } else { "import" };
1106 self.db
1107 .prepare("UPDATE github_repos SET mode = ? WHERE repo_id = ?")
1108 .bind(&[mode.into(), row.repo_id.as_str().into()])?
1109 .run()
1110 .await?;
1111 }
1112 let title = if a.keep_remote_updated {
1113 format!("{} moved {} to g1t; {} now follows it", a.actor.username, row.repo, row.name)
1114 } else {
1115 format!("{} moved {} to g1t and stopped tracking {}", a.actor.username, row.repo, row.name)
1116 };
1117 self.publish("mirror.moved_in", &row, Some(&a.actor.username), title, None, Vec::new()).await;
1118 let mut view = self.view_of(&row.repo_id, true, None).await?;
1119 view.notes = notes;
1120 Ok(Outcome::Ok(view))
1121 }
1122
1123 pub async fn sync_now(&self, a: MirrorActArgs) -> Result<Outcome<MirrorView>> {
1124 let Some(repo) = self.repo(&a.repo_id, Some(&a.actor)).await? else {
1125 return Ok(fail(FailureCode::NotFound, "Repository not found."));
1126 };
1127 // Syncing brings commits in, as pushing does.
1128 if let Some(refused) = Self::refused(&a.actor, &repo, Capability::Push) {
1129 return Ok(refused);
1130 }
1131 let rows = self.rows(&repo.id).await?;
1132 if rows.is_empty() {
1133 return Ok(fail(FailureCode::NotFound, "This repository is not linked to another host."));
1134 }
1135 let mut problems = Vec::new();
1136 for row in &rows {
1137 if row.leads() && !row.state().mirror().is_some_and(MirrorState::follows) {
1138 continue;
1139 }
1140 if let Err(reason) = self.sync(row).await? {
1141 problems.push(format!("{}: {reason}", row.name));
1142 }
1143 }
1144 if !problems.is_empty() {
1145 return Ok(fail(FailureCode::Conflict, problems.join("; ")));
1146 }
1147 Ok(Outcome::Ok(self.view_of(&repo.id, true, None).await?))
1148 }
1149
1150 pub async fn settings(&self, a: MirrorSettingsArgs) -> Result<Outcome<Remote>> {
1151 let Some(row) = self.row(&a.remote_id).await? else {
1152 return Ok(fail(FailureCode::NotFound, "That link was not found."));
1153 };
1154 let Some(repo) = self.repo(&row.repo_id, Some(&a.actor)).await? else {
1155 return Ok(fail(FailureCode::NotFound, "Repository not found."));
1156 };
1157 if let Some(refused) = Self::refused(&a.actor, &repo, Capability::ManageIntegrations) {
1158 return Ok(refused);
1159 }
1160 let (least, most) = TAKE_OVER_AFTER_MINUTES;
1161 if a.settings.take_over_after.is_some_and(|minutes| minutes < least || minutes > most) {
1162 return Ok(fail(FailureCode::Invalid, format!("Take over after between {least} minutes and {} hours.", most / 60)));
1163 }
1164 let settings = serde_json::to_string(&a.settings).unwrap_or_else(|_| "{}".to_owned());
1165 self.db
1166 .prepare("UPDATE remotes SET settings = ?2, recorded = 0 WHERE id = ?1")
1167 .bind(&[row.id.as_str().into(), settings.as_str().into()])?
1168 .run()
1169 .await?;
1170 let row = RemoteRow { settings, ..row };
1171 self.record(&row).await?;
1172 let hosts = self.hosts().await?;
1173 Ok(Outcome::Ok(row.contract(hosts.get(&row.host()))))
1174 }
1175
1176 pub async fn add(&self, a: MirrorAddArgs) -> Result<Outcome<Remote>> {
1177 if a.provider == RemoteProvider::Github {
1178 return Ok(fail(FailureCode::Invalid, "Link GitHub repositories through the GitHub App, from New → Import from GitHub."));
1179 }
1180 let Some(repo) = self.repo(&a.repo_id, Some(&a.actor)).await? else {
1181 return Ok(fail(FailureCode::NotFound, "Repository not found."));
1182 };
1183 if let Some(refused) = Self::refused(&a.actor, &repo, Capability::ManageIntegrations) {
1184 return Ok(refused);
1185 }
1186 let Some(clone_url) = clean_clone_url(&a.url) else {
1187 return Ok(fail(FailureCode::Invalid, "Give the remote's https address, such as https://g1t.sh/acme/web.git."));
1188 };
1189 let Some(token) = a.token.as_deref().map(str::trim).filter(|token| !token.is_empty()) else {
1190 return Ok(fail(FailureCode::Invalid, "Give a token that can read and push to the remote."));
1191 };
1192 let Some(sealer) = &self.sealer else {
1193 return Ok(fail(FailureCode::Conflict, "Tokens cannot be kept on this g1t: INTEGRATIONS_KEY is not set."));
1194 };
1195 if a.role == RemoteRole::Leader && self.leader(&repo.id).await?.is_some() {
1196 return Ok(fail(FailureCode::Conflict, "This repository already mirrors a remote. A repository follows one leader."));
1197 }
1198 let now = now_ms();
1199 let id = new_id("rmt", now);
1200 let row = RemoteRow {
1201 id: id.clone(),
1202 repo_id: repo.id.clone(),
1203 workspace: repo.namespace.clone(),
1204 repo: format!("{}/{}", repo.namespace, repo.name),
1205 provider: a.provider.as_str().to_owned(),
1206 role: a.role.as_str().to_owned(),
1207 name: display_name(&clone_url),
1208 url: web_url(&clone_url),
1209 clone_url: clone_url.clone(),
1210 connection_id: None,
1211 username: a.username.map(|name| name.trim().to_owned()).filter(|name| !name.is_empty()),
1212 credential: Some(sealer.seal(token, &format!("rmt:{id}"))),
1213 state: if a.role == RemoteRole::Leader { "standby" } else { "following" }.to_owned(),
1214 state_since: rfc3339(now),
1215 state_by: Some(a.actor.username.clone()),
1216 settings: "{}".to_owned(),
1217 synced_at: None,
1218 last_error: None,
1219 created_at: rfc3339(now),
1220 };
1221 if row.name.eq_ignore_ascii_case(&format!("g1t.sh/{}", row.repo)) {
1222 return Ok(fail(FailureCode::Invalid, "That is this repository."));
1223 }
1224 // A leader fills an empty repository; one with history of its own
1225 // would lose it.
1226 if a.role == RemoteRole::Leader {
1227 let refs = match self.repos_refs(&row, row.username.clone(), token.to_owned()).await? {
1228 Outcome::Ok(refs) => refs,
1229 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
1230 };
1231 if !refs.ours.is_empty() {
1232 return Ok(fail(
1233 FailureCode::Conflict,
1234 "Only an empty repository can become a mirror: it is filled from the remote. Make a new repository for it.",
1235 ));
1236 }
1237 if let Some(reason) = refs.unreachable {
1238 return Ok(fail(FailureCode::Unavailable, reason));
1239 }
1240 }
1241 self.db
1242 .prepare(
1243 "INSERT INTO remotes
1244 (id, repo_id, workspace, repo, provider, role, name, url, clone_url, username, credential,
1245 state, state_since, state_by, settings, created_by, created_at)
1246 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, '{}', ?15, ?13)",
1247 )
1248 .bind(&[
1249 row.id.as_str().into(),
1250 row.repo_id.as_str().into(),
1251 row.workspace.as_str().into(),
1252 row.repo.as_str().into(),
1253 row.provider.as_str().into(),
1254 row.role.as_str().into(),
1255 row.name.as_str().into(),
1256 row.url.as_str().into(),
1257 row.clone_url.as_str().into(),
1258 null_or(row.username.as_deref()),
1259 null_or(row.credential.as_deref()),
1260 row.state.as_str().into(),
1261 row.state_since.as_str().into(),
1262 a.actor.username.as_str().into(),
1263 a.actor.id.as_str().into(),
1264 ])?
1265 .run()
1266 .await?;
1267 self.record(&row).await?;
1268 let synced = self.sync(&row).await?;
1269 let row = self.row(&id).await?.unwrap_or(row);
1270 if let Err(reason) = synced {
1271 let hosts = self.hosts().await?;
1272 let mut remote = row.contract(hosts.get(&row.host()));
1273 remote.last_error = Some(reason);
1274 return Ok(Outcome::Ok(remote));
1275 }
1276 let hosts = self.hosts().await?;
1277 Ok(Outcome::Ok(row.contract(hosts.get(&row.host()))))
1278 }
1279
1280 pub async fn remove(&self, a: MirrorRemoveArgs) -> Result<Outcome<bool>> {
1281 let Some(row) = self.row(&a.remote_id).await? else {
1282 return Ok(Outcome::Ok(false));
1283 };
1284 let Some(repo) = self.repo(&row.repo_id, Some(&a.actor)).await? else {
1285 return Ok(fail(FailureCode::NotFound, "Repository not found."));
1286 };
1287 if let Some(refused) = Self::refused(&a.actor, &repo, Capability::ManageIntegrations) {
1288 return Ok(refused);
1289 }
1290 if row.leads() && matches!(row.state(), RemoteState::Takeover | RemoteState::HandingBack) {
1291 return Ok(fail(
1292 FailureCode::Conflict,
1293 "g1t leads this repository for now. Hand it back, or move it to g1t, before unlinking.",
1294 ));
1295 }
1296 if row.leads() {
1297 let done: Outcome<Repo> =
1298 g1t_kit::call(&self.repos, "set_mirror", &SetMirrorArgs { repo_id: row.repo_id.clone(), mirror: None }).await?;
1299 if let Outcome::Fail(failure) = done {
1300 return Ok(Outcome::Fail(failure));
1301 }
1302 }
1303 self.db.prepare("DELETE FROM remote_refs WHERE remote_id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1304 self.db.prepare("DELETE FROM remotes WHERE id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1305 if row.provider() == RemoteProvider::Github {
1306 self.db
1307 .prepare("UPDATE github_repos SET mode = 'import' WHERE repo_id = ?")
1308 .bind(&[row.repo_id.as_str().into()])?
1309 .run()
1310 .await?;
1311 }
1312 Ok(Outcome::Ok(true))
1313 }
1314
1315 /// A remote g1t can no longer reach for good (uninstalled, deleted). A
1316 /// mirror standing by becomes an ordinary repository with what it has,
1317 /// since it can no longer follow; a follower is let go; a takeover
1318 /// keeps going and is told why, so nothing is lost.
1319 pub async fn link_gone(&self, id: &str, why: &str) -> Result<()> {
1320 let Some(row) = self.row(id).await? else { return Ok(()) };
1321 if row.leads() && matches!(row.state(), RemoteState::Takeover | RemoteState::HandingBack) {
1322 let reason = format!("{} {why}. g1t keeps leading; move it to g1t to keep it here.", row.name);
1323 return self.note(&row, Some(&reason)).await;
1324 }
1325 if row.leads() {
1326 let _: Outcome<Repo> =
1327 g1t_kit::call(&self.repos, "set_mirror", &SetMirrorArgs { repo_id: row.repo_id.clone(), mirror: None }).await?;
1328 }
1329 self.db.prepare("DELETE FROM remote_refs WHERE remote_id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1330 self.db.prepare("DELETE FROM remotes WHERE id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1331 let title = if row.leads() {
1332 format!("{} stopped mirroring {}: it {why}. g1t keeps what it has.", row.repo, row.name)
1333 } else {
1334 format!("{} stopped pushing to {}: it {why}.", row.repo, row.name)
1335 };
1336 let notify = self.told(&row, None).await;
1337 self.publish("mirror.moved_in", &row, Some("g1t"), title, None, notify).await;
1338 Ok(())
1339 }
1340
1341 // --- What hosts and g1t say -----------------------------------------------
1342
1343 /// A push on GitHub, from the App's webhook.
1344 pub async fn on_github_push(&self, payload: &Value) -> Result<()> {
1345 let Some(id) = payload["repository"]["id"].as_u64() else { return Ok(()) };
1346 let rows = self
1347 .db
1348 .prepare("SELECT * FROM remotes WHERE provider = 'github' AND external_id = ?")
1349 .bind(&[id.to_string().into()])?
1350 .all()
1351 .await?
1352 .results::<RemoteRow>()?;
1353 for row in rows {
1354 if let Err(error) = self.on_remote_push(&row, payload).await {
1355 worker::console_error!("mirrors: a push on {} was not followed: {error}", row.name);
1356 }
1357 }
1358 Ok(())
1359 }
1360
1361 /// A push on a remote: a mirror standing by catches up; during a
1362 /// takeover nothing moves (the hand-back will see it); a follower takes
1363 /// in fast-forwards, or overwrites, as its settings say.
1364 async fn on_remote_push(&self, row: &RemoteRow, payload: &Value) -> Result<()> {
1365 if row.leads() {
1366 if row.state().mirror().is_some_and(MirrorState::follows)
1367 && let Err(reason) = self.sync(row).await?
1368 {
1369 worker::console_log!("mirrors: {} did not catch up with {}: {reason}", row.repo, row.name);
1370 }
1371 return Ok(());
1372 }
1373 let (Some(git_ref), Some(after)) = (payload["ref"].as_str(), payload["after"].as_str()) else {
1374 return Ok(());
1375 };
1376 let before = payload["before"].as_str().filter(|hash| !hash.chars().all(|c| c == '0'));
1377 let deleted = payload["deleted"].as_bool() == Some(true) || after.chars().all(|c| c == '0');
1378 let (username, token) = match self.credential(row).await? {
1379 Ok(credential) => credential,
1380 Err(reason) => return self.note(row, Some(&reason)).await,
1381 };
1382 let refs = match self.repos_refs(row, username.clone(), token.clone()).await? {
1383 Outcome::Ok(refs) => refs,
1384 Outcome::Fail(failure) => return self.note(row, Some(&failure.message)).await,
1385 };
1386 let ours = refs.ours.get(git_ref).map(String::as_str);
1387 // g1t's own push coming back: nothing to do.
1388 if ours == Some(after) || (deleted && ours.is_none()) {
1389 return Ok(());
1390 }
1391 if row.settings().remote_pushes == RemotePushes::Overwrite {
1392 let _ = self.sync(row).await?;
1393 return Ok(());
1394 }
1395 let forced = payload["forced"].as_bool() == Some(true);
1396 if !deleted && !forced && before == ours {
1397 let moves = vec![RefMove {
1398 direction: MirrorDirection::Pull,
1399 git_ref: git_ref.to_owned(),
1400 to: None,
1401 old: ours.map(str::to_owned),
1402 new: Some(after.to_owned()),
1403 }];
1404 if let Outcome::Ok(applied) = self.apply(row, username, token, moves).await?
1405 && applied.moved.iter().all(|moved| moved.problem.is_none())
1406 {
1407 return self.note(row, None).await;
1408 }
1409 }
1410 let branch = branch_of(git_ref).unwrap_or(git_ref);
1411 let reason = format!(
1412 "{branch} changed on {} in a way g1t cannot follow. Sync to push g1t's {branch} over it, or bring that change into g1t first.",
1413 row.name
1414 );
1415 self.note(row, Some(&reason)).await?;
1416 self.set_state(row, RemoteState::Stuck, None).await?;
1417 Ok(())
1418 }
1419
1420 /// A push on g1t: each follower is sent it. A stuck follower waits for
1421 /// someone to sync it, so a change made on it is not pushed over.
1422 pub async fn on_event(&self, event: &g1t_contracts::events::Event) -> Result<()> {
1423 if event.kind != "git.push" {
1424 return Ok(());
1425 }
1426 let Some(repo_id) = event.repo_id.as_deref() else { return Ok(()) };
1427 for row in self.rows(repo_id).await? {
1428 if row.leads() || row.state() != RemoteState::Following {
1429 continue;
1430 }
1431 if let Err(reason) = self.sync(&row).await? {
1432 worker::console_log!("mirrors: {} not pushed to {}: {reason}", row.repo, row.name);
1433 }
1434 }
1435 Ok(())
1436 }
1437
1438 /// Records one check of a host, and acts when it goes down or comes back.
1439 async fn host_checked(&self, host: &str, answered: bool, problem: Option<&str>) -> Result<()> {
1440 let now = now_ms();
1441 let row = self
1442 .db
1443 .prepare("SELECT * FROM remote_hosts WHERE host = ?")
1444 .bind(&[host.into()])?
1445 .first::<HostRow>(None)
1446 .await?
1447 .unwrap_or_else(|| HostRow { host: host.to_owned(), ..HostRow::default() });
1448 let (next, change) = row.health().checked(answered, now);
1449 self.db
1450 .prepare(
1451 "INSERT INTO remote_hosts (host, failures, successes, streak_ms, unreachable_since, checked_ms, last_problem)
1452 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
1453 ON CONFLICT (host) DO UPDATE SET failures = excluded.failures, successes = excluded.successes,
1454 streak_ms = excluded.streak_ms, unreachable_since = excluded.unreachable_since,
1455 checked_ms = excluded.checked_ms, last_problem = COALESCE(excluded.last_problem, remote_hosts.last_problem)",
1456 )
1457 .bind(&[
1458 host.into(),
1459 next.failures.into(),
1460 next.successes.into(),
1461 (next.streak_ms as f64).into(),
1462 null_or(next.unreachable_since.map(rfc3339).as_deref()),
1463 (now as f64).into(),
1464 null_or(problem),
1465 ])?
1466 .run()
1467 .await?;
1468 if change == Change::None {
1469 return Ok(());
1470 }
1471 let rows = self
1472 .db
1473 .prepare("SELECT * FROM remotes WHERE role = 'leader'")
1474 .all()
1475 .await?
1476 .results::<RemoteRow>()?;
1477 for row in rows.into_iter().filter(|row| row.host() == host) {
1478 let (kind, title) = match change {
1479 Change::WentDown => ("mirror.unreachable", format!("{} is not answering. {} keeps its copy.", row.name, row.repo)),
1480 _ => ("mirror.reachable", format!("{} answers again", row.name)),
1481 };
1482 let notify = self.told(&row, None).await;
1483 let detail = (change == Change::WentDown && row.state().mirror().is_some_and(MirrorState::follows))
1484 .then(|| format!("Take over {} to keep working on g1t until it is back.", row.repo));
1485 self.publish(kind, &row, None, title, detail, notify).await;
1486 }
1487 Ok(())
1488 }
1489
1490 // --- Every minute -----------------------------------------------------------
1491
1492 pub async fn on_minute(&self) -> Result<()> {
1493 let now = now_ms();
1494 // Mirrors the repos service has not heard of yet.
1495 let unrecorded = self
1496 .db
1497 .prepare("SELECT * FROM remotes WHERE recorded = 0 LIMIT 50")
1498 .all()
1499 .await?
1500 .results::<RemoteRow>()?;
1501 for row in &unrecorded {
1502 if row.leads() {
1503 self.record(row).await?;
1504 } else {
1505 self.db.prepare("UPDATE remotes SET recorded = 1 WHERE id = ?").bind(&[row.id.as_str().into()])?.run().await?;
1506 }
1507 }
1508 let leaders = self
1509 .db
1510 .prepare("SELECT * FROM remotes WHERE role = 'leader'")
1511 .all()
1512 .await?
1513 .results::<RemoteRow>()?;
1514 // One check per host a mirror follows.
1515 let mut by_host: BTreeMap<String, &RemoteRow> = BTreeMap::new();
1516 for row in &leaders {
1517 by_host.entry(row.host()).or_insert(row);
1518 }
1519 for (host, row) in &by_host {
1520 let answered = self.probe(row).await?;
1521 self.host_checked(host, answered.is_ok(), answered.err().as_deref()).await?;
1522 }
1523 let hosts = self.hosts().await?;
1524 for row in &leaders {
1525 let health = hosts.get(&row.host()).map(HostRow::health).unwrap_or_default();
1526 let settings = row.settings();
1527 match row.state() {
1528 // Taking over on its own, only when asked to.
1529 RemoteState::Standby | RemoteState::Ci => {
1530 if let (Some(minutes), Some(since)) = (settings.take_over_after, health.unreachable_since)
1531 && now.saturating_sub(since) >= u64::from(minutes) * 60_000
1532 {
1533 if let Outcome::Fail(failure) = self.start_takeover(row.clone(), By::G1t).await? {
1534 worker::console_log!("mirrors: {} not taken over: {}", row.repo, failure.message);
1535 }
1536 continue;
1537 }
1538 // Hosts with no webhook are polled.
1539 if row.provider() != RemoteProvider::Github && health.reachable() {
1540 let polled = self
1541 .db
1542 .prepare("UPDATE remotes SET polled_ms = ?2 WHERE id = ?1 AND polled_ms < ?3 RETURNING id")
1543 .bind(&[row.id.as_str().into(), (now as f64).into(), (now.saturating_sub(POLL_MS) as f64).into()])?
1544 .first::<Value>(None)
1545 .await?;
1546 if polled.is_some() {
1547 let _ = self.sync(row).await?;
1548 }
1549 }
1550 }
1551 // Handing back on its own once the remote is back and it is
1552 // clean, when asked to.
1553 RemoteState::Takeover if health.reachable() && settings.hand_back == HandBack::WhenClean && row.state_by.as_deref() == Some("g1t") => {
1554 if let Outcome::Ok(plan) = self.plan(row).await?
1555 && plan.clean()
1556 && let Outcome::Fail(failure) = self.finish_takeover(row.clone(), By::G1t).await?
1557 {
1558 worker::console_log!("mirrors: {} not handed back: {}", row.repo, failure.message);
1559 }
1560 }
1561 _ => {}
1562 }
1563 }
1564 Ok(())
1565 }
1566
1567 /// Whether a remote's host answers `info/refs`: any answer but a server
1568 /// error counts, since a refused credential is the link's problem, not
1569 /// the host's.
1570 async fn probe(&self, row: &RemoteRow) -> Result<std::result::Result<(), String>> {
1571 let (username, token) = match self.credential(row).await? {
1572 Ok(credential) => credential,
1573 // No credential: ask without one; the host still answers.
1574 Err(_) => (None, String::new()),
1575 };
1576 let authorization = if token.is_empty() {
1577 String::new()
1578 } else {
1579 let user = username.unwrap_or_else(|| "x-access-token".to_owned());
1580 format!("Basic {}", base64(&format!("{user}:{token}")))
1581 };
1582 let url = format!("{}/info/refs?service=git-upload-pack", row.clone_url);
1583 let mut headers = vec![("accept", "*/*")];
1584 if !authorization.is_empty() {
1585 headers.push(("authorization", authorization.as_str()));
1586 }
1587 match http::send(Method::Get, &url, &headers, None).await {
1588 Ok(answer) if answer.status < 500 && answer.status != 429 && answer.status != 408 => Ok(Ok(())),
1589 Ok(answer) => Ok(Err(format!("{} answered {}.", row.host(), answer.status))),
1590 Err(_) => Ok(Err(format!("{} could not be reached.", row.host()))),
1591 }
1592 }
1593}
1594
1595fn base64(text: &str) -> String {
1596 const ALPHABET: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
1597 let bytes = text.as_bytes();
1598 let mut out = String::with_capacity(bytes.len().div_ceil(3) * 4);
1599 for chunk in bytes.chunks(3) {
1600 let n = (chunk[0] as u32) << 16 | (*chunk.get(1).unwrap_or(&0) as u32) << 8 | *chunk.get(2).unwrap_or(&0) as u32;
1601 for (i, shift) in [18, 12, 6, 0].into_iter().enumerate() {
1602 out.push(if i <= chunk.len() { ALPHABET[(n >> shift) as usize & 63] as char } else { '=' });
1603 }
1604 }
1605 out
1606}
1607
1608/// When the remote is away, the plan assumes it is where the takeover
1609/// started.
1610fn bases_as_theirs(bases: &BTreeMap<String, (Option<String>, Option<RefDecision>)>) -> BTreeMap<String, String> {
1611 bases.iter().filter_map(|(name, (base, _))| base.clone().map(|base| (name.clone(), base))).collect()
1612}
1613
1614fn undecided(plan: &HandbackPlan) -> g1t_contracts::Failure {
1615 let waiting: Vec<&str> = plan
1616 .refs
1617 .iter()
1618 .filter(|item| item.action == RefAction::Diverged && item.decision.is_none())
1619 .map(|item| branch_of(&item.git_ref).unwrap_or(&item.git_ref))
1620 .collect();
1621 g1t_contracts::Failure {
1622 code: FailureCode::Conflict,
1623 message: format!(
1624 "Changed on both sides: {}. Decide for each whether to keep g1t's, keep the remote's, or send g1t's as a pull request.",
1625 waiting.join(", ")
1626 ),
1627 }
1628}
1629
1630/// Answers the `mirror_*` methods; `None` for any other.
1631pub async fn route(method: &str, body: &Value, env: &Env, _ctx: &Context) -> Option<Result<Response>> {
1632 if !method.starts_with("mirror_") {
1633 return None;
1634 }
1635 Some(handle(method, body.clone(), env).await)
1636}
1637
1638async fn handle(method: &str, body: Value, env: &Env) -> Result<Response> {
1639 let mirrors = Mirrors::new(env)?;
1640 match method {
1641 "mirror_view" => reply(&mirrors.view(args(body)?).await?),
1642 "mirror_briefs" => reply(&mirrors.briefs(args(body)?).await?),
1643 "mirror_take_over" => reply(&mirrors.take_over(args(body)?).await?),
1644 "mirror_ci" => reply(&mirrors.ci(args(body)?).await?),
1645 "mirror_hand_back_plan" => reply(&mirrors.hand_back_plan(args(body)?).await?),
1646 "mirror_hand_back" => reply(&mirrors.hand_back(args(body)?).await?),
1647 "mirror_move_in" => reply(&mirrors.move_in(args(body)?).await?),
1648 "mirror_sync" => reply(&mirrors.sync_now(args(body)?).await?),
1649 "mirror_settings" => reply(&mirrors.settings(args(body)?).await?),
1650 "mirror_add" => reply(&mirrors.add(args(body)?).await?),
1651 "mirror_remove" => reply(&mirrors.remove(args(body)?).await?),
1652 _ => Response::error("Unknown method", 404),
1653 }
1654}
1655
1656#[cfg(test)]
1657mod tests {
1658 use super::*;
1659
1660 #[test]
1661 fn remotes_are_named_as_people_say_them() {
1662 assert_eq!(host_of("https://GitHub.com/acme/web.git"), "github.com");
1663 assert_eq!(host_of("https://user@git.acme.internal/x"), "git.acme.internal");
1664 assert_eq!(display_name("https://github.com/acme/web.git"), "github.com/acme/web");
1665 assert_eq!(display_name("https://g1t.sh/acme/web/"), "g1t.sh/acme/web");
1666 assert_eq!(clean_clone_url("https://g1t.sh/acme/web").as_deref(), Some("https://g1t.sh/acme/web.git"));
1667 assert_eq!(clean_clone_url("https://git.example/a.git/").as_deref(), Some("https://git.example/a.git"));
1668 assert_eq!(clean_clone_url("http://git.example/a"), None, "https only");
1669 assert_eq!(clean_clone_url("https://me:secret@git.example/a"), None, "no credentials in the address");
1670 assert_eq!(clean_clone_url("https://git.example"), None);
1671 assert_eq!(web_url("https://github.com/acme/web.git"), "https://github.com/acme/web");
1672 }
1673
1674 #[test]
1675 fn a_host_goes_down_slowly_and_comes_back_slower() {
1676 let minute = 60_000;
1677 let mut health = HostHealth::default();
1678 let mut changes = Vec::new();
1679 for at in [0, minute, 2 * minute] {
1680 let (next, change) = health.checked(false, at);
1681 health = next;
1682 changes.push(change);
1683 }
1684 assert_eq!(changes, [Change::None, Change::None, Change::WentDown], "three failures over two minutes");
1685 assert!(!health.reachable());
1686 // One blip does not bring it back, nor do three quick answers.
1687 let (next, change) = health.checked(true, 3 * minute);
1688 assert_eq!(change, Change::None);
1689 let (next, _) = next.checked(true, 3 * minute + 1);
1690 let (next, change) = next.checked(true, 3 * minute + 2);
1691 assert_eq!(change, Change::None, "three answers in a moment are not five minutes");
1692 let (back, change) = next.checked(true, 8 * minute);
1693 assert_eq!(change, Change::CameBack);
1694 assert!(back.reachable());
1695 // A failure in between starts the count again.
1696 let (flaky, _) = HostHealth::default().checked(false, 0);
1697 let (flaky, _) = flaky.checked(true, minute);
1698 let (flaky, _) = flaky.checked(false, 2 * minute);
1699 let (_, change) = flaky.checked(false, 3 * minute);
1700 assert_eq!(change, Change::None);
1701 }
1702
1703 fn refs(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
1704 pairs.iter().map(|(name, hash)| (name.to_string(), hash.to_string())).collect()
1705 }
1706
1707 #[test]
1708 fn a_hand_back_plan_reads_each_ref_against_where_the_takeover_began() {
1709 let bases: BTreeMap<_, _> = [
1710 ("refs/heads/main", "a"),
1711 ("refs/heads/fix", "f"),
1712 ("refs/heads/docs", "d"),
1713 ("refs/heads/quiet", "q"),
1714 ]
1715 .into_iter()
1716 .map(|(name, base)| (name.to_owned(), (Some(base.to_owned()), None)))
1717 .collect();
1718 let ours = refs(&[
1719 ("refs/heads/main", "a2"),
1720 ("refs/heads/fix", "f2"),
1721 ("refs/heads/docs", "d2"),
1722 ("refs/heads/quiet", "q"),
1723 ("refs/heads/agent", "n"),
1724 ]);
1725 let theirs = refs(&[
1726 ("refs/heads/main", "a"),
1727 ("refs/heads/fix", "f"),
1728 ("refs/heads/docs", "d3"),
1729 ("refs/heads/quiet", "q"),
1730 ("refs/heads/g1t/handback/old", "x"),
1731 ]);
1732 let protected = BTreeSet::from(["refs/heads/main".to_owned()]);
1733 let plan = plan_refs(&bases, &ours, &theirs, &protected);
1734 let actions: Vec<(&str, RefAction)> = plan.iter().map(|item| (item.git_ref.as_str(), item.action)).collect();
1735 assert_eq!(actions, [
1736 ("refs/heads/agent", RefAction::Push),
1737 ("refs/heads/docs", RefAction::Diverged),
1738 ("refs/heads/fix", RefAction::Push),
1739 ("refs/heads/main", RefAction::PullRequest),
1740 ]);
1741 let plan = HandbackPlan::new(plan, true);
1742 assert!(!plan.ready, "docs needs a decision");
1743
1744 let (moves, pull_requests) = hand_back_moves(&plan, &theirs);
1745 assert_eq!(pull_requests, ["refs/heads/main"]);
1746 let summary: Vec<(MirrorDirection, &str, Option<&str>)> =
1747 moves.iter().map(|m| (m.direction, m.git_ref.as_str(), m.to.as_deref())).collect();
1748 assert_eq!(summary, [
1749 (MirrorDirection::Push, "refs/heads/agent", None),
1750 (MirrorDirection::Push, "refs/heads/fix", None),
1751 (MirrorDirection::Push, "refs/heads/main", Some("refs/heads/g1t/handback/main")),
1752 (MirrorDirection::Pull, "refs/heads/main", None),
1753 ], "an undecided ref does not move");
1754 let main_back = &moves[3];
1755 assert_eq!((main_back.old.as_deref(), main_back.new.as_deref()), (Some("a2"), Some("a")), "g1t follows the remote's main");
1756 let push_fix = &moves[1];
1757 assert_eq!((push_fix.old.as_deref(), push_fix.new.as_deref()), (Some("f"), Some("f2")), "only from where the remote is");
1758 }
1759
1760 #[test]
1761 fn decisions_settle_a_diverged_ref() {
1762 let plan = |decision| {
1763 HandbackPlan::new(
1764 vec![RefPlan {
1765 git_ref: "refs/heads/docs".into(),
1766 base: Some("d".into()),
1767 ours: Some("d2".into()),
1768 theirs: Some("d3".into()),
1769 action: RefAction::Diverged,
1770 decision: Some(decision),
1771 }],
1772 true,
1773 )
1774 };
1775 let theirs = refs(&[("refs/heads/g1t/handback/docs", "old")]);
1776 let (moves, _) = hand_back_moves(&plan(RefDecision::KeepOurs), &theirs);
1777 assert_eq!((moves[0].direction, moves[0].old.as_deref(), moves[0].new.as_deref()), (MirrorDirection::Push, Some("d3"), Some("d2")));
1778 let (moves, _) = hand_back_moves(&plan(RefDecision::KeepTheirs), &theirs);
1779 assert_eq!((moves[0].direction, moves[0].old.as_deref(), moves[0].new.as_deref()), (MirrorDirection::Pull, Some("d2"), Some("d3")));
1780 let (moves, pull_requests) = hand_back_moves(&plan(RefDecision::PullRequest), &theirs);
1781 assert_eq!(pull_requests, ["refs/heads/docs"]);
1782 assert_eq!(moves[0].old.as_deref(), Some("old"), "an old hand-back branch is moved from where it is");
1783 assert_eq!(moves.len(), 2);
1784 }
1785}