| 1 | //! Moving a repository from one git store namespace to another |
| 2 | //! (docs/ARTIFACTS.md, R7), keeping its id, its path and everything g1t |
| 3 | //! knows about it. Only its store key changes. |
| 4 | //! |
| 5 | //! A move is asked for (`move_repository`, or a row an operator inserts |
| 6 | //! with `scripts/ops/artifacts-namespaces.mjs move`) and run by the hourly |
| 7 | //! sweep, one at a time: |
| 8 | //! |
| 9 | //! 1. **Paused.** Writes to the repository and its pull requests' working |
| 10 | //! copies wait (`writes_paused_until`, at most [`PAUSE_MS`]): pushes wait |
| 11 | //! up to 20 seconds and are then told to try again, as are merges, |
| 12 | //! commits from the web and push credentials for sandboxes. A move |
| 13 | //! waits for push credentials already handed out to expire (they reach |
| 14 | //! the store directly), then a few seconds for pushes in flight. |
| 15 | //! 2. **Copied.** Each one is made in the new namespace under the same |
| 16 | //! name, and every ref is copied over git: one upload-pack from the old |
| 17 | //! copy streamed into one receive-pack to the new (land.rs `copy_refs`), |
| 18 | //! so its size is bounded by time, not memory. Until both copies list the |
| 19 | //! same refs and the old one has not moved since, it is copied again, |
| 20 | //! three times at most. |
| 21 | //! 3. **Switched.** Every row's `store` names the new key in one batch, and |
| 22 | //! its `refs_version` moves, so nothing kept for the old copy is used |
| 23 | //! again. Writes go on. Removed working copies (forks.rs) need no |
| 24 | //! copying: their rows follow, so one made again is made in the new |
| 25 | //! namespace. |
| 26 | //! 4. **Cleaned.** After [`MOVE_KEEP_DAYS`] the old copies are deleted, each |
| 27 | //! only while its refs still say what was copied. One that changed (a |
| 28 | //! push that slipped past the pause) is kept, and the move marked |
| 29 | //! `diverged` for an operator. |
| 30 | //! |
| 31 | //! A failure before the switch deletes what was made in the new namespace |
| 32 | //! and lifts the pause; the repository stays where it was. |
| 33 | |
| 34 | use std::collections::BTreeMap; |
| 35 | |
| 36 | use g1t_contracts::repos::Repo; |
| 37 | use g1t_contracts::{FailureCode, Outcome, new_id}; |
| 38 | use g1t_kit::now_ms; |
| 39 | use serde::{Deserialize, Serialize}; |
| 40 | use worker::wasm_bindgen::JsValue; |
| 41 | use worker::Result; |
| 42 | |
| 43 | use crate::registry::{Registry, store_key}; |
| 44 | use crate::store::{GitRepo, GitStore, Scope, locate}; |
| 45 | use crate::{Repos, land, mirror, refs, shards}; |
| 46 | |
| 47 | /// The longest writes wait for one move. |
| 48 | pub const PAUSE_MS: u64 = 20 * 60 * 1000; |
| 49 | /// How long a write waits for a pause to end before it is told to try again. |
| 50 | pub const PAUSE_WAIT_MS: u64 = 20_000; |
| 51 | const PAUSE_POLL_MS: u64 = 2_000; |
| 52 | /// How long a move waits for push credentials already handed out; longer |
| 53 | /// and it is tried again in the next sweep. |
| 54 | const OPEN_WAIT_MS: u64 = 7 * 60 * 1000; |
| 55 | /// Pushes already past the pause get this long to finish. |
| 56 | const SETTLE_MS: u64 = 5_000; |
| 57 | /// Rounds of copying before a move that keeps changing gives up. |
| 58 | const ROUNDS: u32 = 3; |
| 59 | /// Days the old copies are kept after a move. |
| 60 | pub const MOVE_KEEP_DAYS: u64 = 7; |
| 61 | const DAY_MS: u64 = 24 * 3600 * 1000; |
| 62 | /// Moves run per sweep, and old copies cleaned. |
| 63 | const MOVES_PER_SWEEP: u32 = 1; |
| 64 | const CLEANS_PER_SWEEP: u32 = 10; |
| 65 | /// Tries before a move that keeps failing is left failed. |
| 66 | const MAX_ATTEMPTS: u32 = 3; |
| 67 | const ZERO_ID: &str = "0000000000000000000000000000000000000000"; |
| 68 | |
| 69 | /// `move_repository`: services and operators only. |
| 70 | #[derive(Debug, Deserialize)] |
| 71 | pub struct MoveArgs { |
| 72 | pub repo_id: String, |
| 73 | /// The namespace it goes to. |
| 74 | pub namespace: String, |
| 75 | #[serde(default)] |
| 76 | pub requested_by: Option<String>, |
| 77 | } |
| 78 | |
| 79 | /// One move, as `repository_moves` answers. |
| 80 | #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] |
| 81 | pub struct MoveRow { |
| 82 | pub id: String, |
| 83 | pub repo_id: String, |
| 84 | pub to_namespace: String, |
| 85 | pub status: String, |
| 86 | #[serde(default)] |
| 87 | pub requested_by: Option<String>, |
| 88 | pub queued_ms: f64, |
| 89 | #[serde(default)] |
| 90 | pub started_ms: Option<f64>, |
| 91 | #[serde(default)] |
| 92 | pub finished_ms: Option<f64>, |
| 93 | #[serde(default)] |
| 94 | pub cleaned_ms: Option<f64>, |
| 95 | #[serde(default)] |
| 96 | pub attempts: f64, |
| 97 | #[serde(default)] |
| 98 | pub note: Option<String>, |
| 99 | } |
| 100 | |
| 101 | #[derive(Debug, Deserialize)] |
| 102 | pub struct ListMovesArgs { |
| 103 | #[serde(default)] |
| 104 | pub limit: Option<u32>, |
| 105 | } |
| 106 | |
| 107 | /// What one copy holds, as copied. |
| 108 | #[derive(Debug, Deserialize)] |
| 109 | struct CopyRow { |
| 110 | repo_id: String, |
| 111 | from_key: String, |
| 112 | refs: String, |
| 113 | } |
| 114 | |
| 115 | /// Every ref worth copying, by name: peeled tag lines are the |
| 116 | /// advertisement's, not refs. |
| 117 | pub fn ref_map(all: Vec<(String, String)>) -> BTreeMap<String, String> { |
| 118 | all.into_iter().filter(|(name, _)| !name.ends_with("^{}")).collect() |
| 119 | } |
| 120 | |
| 121 | /// What makes `target`'s refs `source`'s: the commands, the objects to |
| 122 | /// fetch, and what the target holds already. |
| 123 | pub fn plan(source: &BTreeMap<String, String>, target: &BTreeMap<String, String>) -> (Vec<mirror::Command>, Vec<String>, Vec<String>) { |
| 124 | let mut commands: Vec<mirror::Command> = Vec::new(); |
| 125 | for (name, hash) in source { |
| 126 | if target.get(name) != Some(hash) { |
| 127 | commands.push((name.clone(), target.get(name).cloned(), hash.clone())); |
| 128 | } |
| 129 | } |
| 130 | for (name, hash) in target { |
| 131 | if !source.contains_key(name) { |
| 132 | commands.push((name.clone(), Some(hash.clone()), ZERO_ID.to_owned())); |
| 133 | } |
| 134 | } |
| 135 | let held: std::collections::BTreeSet<&String> = target.values().collect(); |
| 136 | let mut wants: Vec<String> = commands |
| 137 | .iter() |
| 138 | .map(|(_, _, new)| new.clone()) |
| 139 | .filter(|new| new != ZERO_ID && !held.contains(new)) |
| 140 | .collect(); |
| 141 | wants.sort(); |
| 142 | wants.dedup(); |
| 143 | let mut haves: Vec<String> = held.into_iter().cloned().collect(); |
| 144 | haves.dedup(); |
| 145 | (commands, wants, haves) |
| 146 | } |
| 147 | |
| 148 | /// What a writer is told while writes wait. |
| 149 | pub fn paused_message(repo: &Repo, reason: &str) -> String { |
| 150 | format!( |
| 151 | "{}/{} is paused for maintenance ({reason}); changes to it wait a few minutes. Try again shortly.", |
| 152 | repo.namespace, repo.name |
| 153 | ) |
| 154 | } |
| 155 | |
| 156 | /// Why a move cannot start, if it cannot. |
| 157 | pub fn refusal(repo: Option<&Repo>, to: &str, bound: &[String], writable: bool, from: &str) -> Option<String> { |
| 158 | let Some(repo) = repo else { |
| 159 | return Some("the repository is gone".to_owned()); |
| 160 | }; |
| 161 | if repo.fork_of.is_some() { |
| 162 | return Some("a pull request's working copy moves with its repository".to_owned()); |
| 163 | } |
| 164 | if !shards::valid_namespace(to) || !bound.iter().any(|name| name == to) { |
| 165 | return Some(format!("{to} is not a bound namespace")); |
| 166 | } |
| 167 | if !writable { |
| 168 | return Some(format!("{to} does not take writes now")); |
| 169 | } |
| 170 | if from == to { |
| 171 | return Some(format!("it is in {to} already")); |
| 172 | } |
| 173 | None |
| 174 | } |
| 175 | |
| 176 | impl Registry { |
| 177 | pub async fn queue_move(&self, a: &MoveArgs, now: u64) -> Result<MoveRow> { |
| 178 | let id = new_id("mov", now); |
| 179 | self.db |
| 180 | .prepare( |
| 181 | "INSERT INTO repo_moves (id, repo_id, to_namespace, status, requested_by, queued_ms) |
| 182 | VALUES (?1, ?2, ?3, 'queued', ?4, ?5)", |
| 183 | ) |
| 184 | .bind(&[ |
| 185 | id.as_str().into(), |
| 186 | a.repo_id.as_str().into(), |
| 187 | a.namespace.as_str().into(), |
| 188 | a.requested_by.as_deref().map_or(JsValue::NULL, JsValue::from), |
| 189 | (now as f64).into(), |
| 190 | ])? |
| 191 | .run() |
| 192 | .await?; |
| 193 | self.move_by_id(&id).await?.ok_or_else(|| worker::Error::RustError("the move was not recorded".to_owned())) |
| 194 | } |
| 195 | |
| 196 | async fn move_by_id(&self, id: &str) -> Result<Option<MoveRow>> { |
| 197 | self.db.prepare("SELECT * FROM repo_moves WHERE id = ?").bind(&[id.into()])?.first::<MoveRow>(None).await |
| 198 | } |
| 199 | |
| 200 | pub async fn moves(&self, limit: u32) -> Result<Vec<MoveRow>> { |
| 201 | self.db |
| 202 | .prepare("SELECT * FROM repo_moves ORDER BY queued_ms DESC LIMIT ?") |
| 203 | .bind(&[limit.clamp(1, 200).into()])? |
| 204 | .all() |
| 205 | .await? |
| 206 | .results::<MoveRow>() |
| 207 | } |
| 208 | |
| 209 | async fn queued_moves(&self, limit: u32) -> Result<Vec<MoveRow>> { |
| 210 | self.db |
| 211 | .prepare("SELECT * FROM repo_moves WHERE status = 'queued' ORDER BY queued_ms LIMIT ?") |
| 212 | .bind(&[limit.into()])? |
| 213 | .all() |
| 214 | .await? |
| 215 | .results::<MoveRow>() |
| 216 | } |
| 217 | |
| 218 | async fn set_move(&self, id: &str, status: &str, note: Option<&str>, now: u64) -> Result<()> { |
| 219 | let finished = matches!(status, "moved" | "failed"); |
| 220 | let cleaned = status == "cleaned"; |
| 221 | self.db |
| 222 | .prepare( |
| 223 | "UPDATE repo_moves SET status = ?2, note = ?3, |
| 224 | started_ms = CASE WHEN ?2 = 'moving' THEN ?4 ELSE started_ms END, |
| 225 | attempts = attempts + CASE WHEN ?2 = 'moving' THEN 1 ELSE 0 END, |
| 226 | finished_ms = CASE WHEN ?5 THEN ?4 ELSE finished_ms END, |
| 227 | cleaned_ms = CASE WHEN ?6 THEN ?4 ELSE cleaned_ms END |
| 228 | WHERE id = ?1", |
| 229 | ) |
| 230 | .bind(&[id.into(), status.into(), note.map_or(JsValue::NULL, JsValue::from), (now as f64).into(), finished.into(), cleaned.into()])? |
| 231 | .run() |
| 232 | .await?; |
| 233 | Ok(()) |
| 234 | } |
| 235 | |
| 236 | /// Writes to these repositories wait until `until`. |
| 237 | async fn pause(&self, ids: &[String], until: u64, reason: &str) -> Result<()> { |
| 238 | let statements = ids |
| 239 | .iter() |
| 240 | .map(|id| { |
| 241 | self.db |
| 242 | .prepare("UPDATE repos SET writes_paused_until = ?, writes_paused_for = ? WHERE id = ?") |
| 243 | .bind(&[(until as f64).into(), reason.into(), id.as_str().into()]) |
| 244 | }) |
| 245 | .collect::<Result<Vec<_>>>()?; |
| 246 | self.db.batch(statements).await?; |
| 247 | for id in ids { |
| 248 | crate::registry::note_paused(id, Some(until), Some(reason)); |
| 249 | } |
| 250 | Ok(()) |
| 251 | } |
| 252 | |
| 253 | async fn resume(&self, ids: &[String]) -> Result<()> { |
| 254 | let statements = ids |
| 255 | .iter() |
| 256 | .map(|id| { |
| 257 | self.db |
| 258 | .prepare("UPDATE repos SET writes_paused_until = NULL, writes_paused_for = NULL WHERE id = ?") |
| 259 | .bind(&[id.as_str().into()]) |
| 260 | }) |
| 261 | .collect::<Result<Vec<_>>>()?; |
| 262 | self.db.batch(statements).await?; |
| 263 | for id in ids { |
| 264 | crate::registry::note_paused(id, None, None); |
| 265 | } |
| 266 | Ok(()) |
| 267 | } |
| 268 | |
| 269 | /// A repository's pull requests' working copies, removed ones too. |
| 270 | async fn all_forks_of(&self, id: &str) -> Result<Vec<Repo>> { |
| 271 | let rows = self |
| 272 | .db |
| 273 | .prepare("SELECT * FROM repos WHERE fork_of = ? AND deleted_at IS NULL") |
| 274 | .bind(&[id.into()])? |
| 275 | .all() |
| 276 | .await? |
| 277 | .results::<crate::registry::RepoRow>()?; |
| 278 | Ok(rows.into_iter().map(Repo::from).collect()) |
| 279 | } |
| 280 | |
| 281 | /// The switch: every row names its new key, its refs version moves, |
| 282 | /// writes go on, and what was copied is recorded, in one batch. |
| 283 | async fn switch(&self, move_id: &str, switched: &[(Repo, String, String, BTreeMap<String, String>)], now: u64) -> Result<()> { |
| 284 | let mut statements = Vec::new(); |
| 285 | for (repo, from, to, refs) in switched { |
| 286 | statements.push( |
| 287 | self.db |
| 288 | .prepare( |
| 289 | "UPDATE repos SET store = ?2, refs_version = coalesce(refs_version, 0) + 1, |
| 290 | writes_paused_until = NULL, writes_paused_for = NULL |
| 291 | WHERE id = ?1", |
| 292 | ) |
| 293 | .bind(&[repo.id.as_str().into(), to.as_str().into()])?, |
| 294 | ); |
| 295 | let (_, name) = crate::shards::split(from); |
| 296 | statements.push( |
| 297 | self.db |
| 298 | .prepare( |
| 299 | "INSERT OR REPLACE INTO repo_move_copies (move_id, repo_id, from_key, to_key, name, refs) |
| 300 | VALUES (?1, ?2, ?3, ?4, ?5, ?6)", |
| 301 | ) |
| 302 | .bind(&[ |
| 303 | move_id.into(), |
| 304 | repo.id.as_str().into(), |
| 305 | from.as_str().into(), |
| 306 | to.as_str().into(), |
| 307 | name.into(), |
| 308 | serde_json::to_string(refs)?.into(), |
| 309 | ])?, |
| 310 | ); |
| 311 | } |
| 312 | statements.push( |
| 313 | self.db |
| 314 | .prepare("UPDATE repo_moves SET status = 'moved', finished_ms = ?2, note = NULL WHERE id = ?1") |
| 315 | .bind(&[move_id.into(), (now as f64).into()])?, |
| 316 | ); |
| 317 | self.db.batch(statements).await?; |
| 318 | for (repo, _, to, _) in switched { |
| 319 | crate::registry::remember_store(repo, to); |
| 320 | crate::registry::note_paused(&repo.id, None, None); |
| 321 | } |
| 322 | Ok(()) |
| 323 | } |
| 324 | |
| 325 | async fn copies_of(&self, move_id: &str) -> Result<Vec<CopyRow>> { |
| 326 | self.db |
| 327 | .prepare("SELECT repo_id, from_key, refs FROM repo_move_copies WHERE move_id = ? AND cleaned_ms IS NULL") |
| 328 | .bind(&[move_id.into()])? |
| 329 | .all() |
| 330 | .await? |
| 331 | .results::<CopyRow>() |
| 332 | } |
| 333 | |
| 334 | async fn copy_cleaned(&self, move_id: &str, repo_id: &str, now: u64) -> Result<()> { |
| 335 | self.db |
| 336 | .prepare("UPDATE repo_move_copies SET cleaned_ms = ?3 WHERE move_id = ?1 AND repo_id = ?2") |
| 337 | .bind(&[move_id.into(), repo_id.into(), (now as f64).into()])? |
| 338 | .run() |
| 339 | .await?; |
| 340 | Ok(()) |
| 341 | } |
| 342 | |
| 343 | async fn moved_before(&self, before: u64, limit: u32) -> Result<Vec<MoveRow>> { |
| 344 | self.db |
| 345 | .prepare("SELECT * FROM repo_moves WHERE status = 'moved' AND finished_ms < ? ORDER BY finished_ms LIMIT ?") |
| 346 | .bind(&[(before as f64).into(), limit.into()])? |
| 347 | .all() |
| 348 | .await? |
| 349 | .results::<MoveRow>() |
| 350 | } |
| 351 | } |
| 352 | |
| 353 | /// How a move went. |
| 354 | enum Ran { |
| 355 | Moved, |
| 356 | /// Not yet: back in the queue, with why. |
| 357 | Later(String), |
| 358 | Failed(String), |
| 359 | } |
| 360 | |
| 361 | impl<S: GitStore> Repos<S> { |
| 362 | /// `move_repository`: queues a move, checked now so a mistake is said |
| 363 | /// at once. |
| 364 | pub(crate) async fn move_repository(&self, a: MoveArgs) -> Result<Outcome<MoveRow>> { |
| 365 | let repo = self.registry.by_id(&a.repo_id).await?; |
| 366 | let from = repo.as_ref().map(|repo| locate(&store_key(repo)).0).unwrap_or_default(); |
| 367 | let bound = self.store.namespaces(); |
| 368 | if let Some(why) = refusal(repo.as_ref(), &a.namespace, &bound, self.store.writable(&a.namespace), &from) { |
| 369 | return Ok(Outcome::fail(FailureCode::Invalid, format!("It cannot be moved: {why}."))); |
| 370 | } |
| 371 | match self.registry.queue_move(&a, now_ms()).await { |
| 372 | Ok(row) => Ok(Outcome::Ok(row)), |
| 373 | // The partial unique index: one in hand already. |
| 374 | Err(error) if error.to_string().contains("UNIQUE") => { |
| 375 | Ok(Outcome::fail(FailureCode::Conflict, "A move of this repository is already queued or running.")) |
| 376 | } |
| 377 | Err(error) => Err(error), |
| 378 | } |
| 379 | } |
| 380 | |
| 381 | /// The sweep's part: runs the queued moves, then cleans old copies. |
| 382 | pub(crate) async fn run_moves(&self) -> Result<u32> { |
| 383 | let mut moved = 0; |
| 384 | for row in self.registry.queued_moves(MOVES_PER_SWEEP).await? { |
| 385 | let now = now_ms(); |
| 386 | self.registry.set_move(&row.id, "moving", None, now).await?; |
| 387 | let ran = match self.run_move(&row).await { |
| 388 | Ok(ran) => ran, |
| 389 | Err(error) => Ran::Failed(error.to_string()), |
| 390 | }; |
| 391 | match ran { |
| 392 | Ran::Moved => moved += 1, |
| 393 | Ran::Later(why) => self.registry.set_move(&row.id, "queued", Some(&why), now_ms()).await?, |
| 394 | Ran::Failed(why) if (row.attempts as u32) + 1 < MAX_ATTEMPTS => { |
| 395 | worker::console_error!("move {} of {} failed, to be tried again: {why}", row.id, row.repo_id); |
| 396 | self.registry.set_move(&row.id, "queued", Some(&why), now_ms()).await?; |
| 397 | } |
| 398 | Ran::Failed(why) => { |
| 399 | worker::console_error!("move {} of {} failed: {why}", row.id, row.repo_id); |
| 400 | self.registry.set_move(&row.id, "failed", Some(&why), now_ms()).await?; |
| 401 | } |
| 402 | } |
| 403 | } |
| 404 | self.clean_moves().await?; |
| 405 | Ok(moved) |
| 406 | } |
| 407 | |
| 408 | async fn run_move(&self, row: &MoveRow) -> Result<Ran> { |
| 409 | let repo = self.registry.by_id(&row.repo_id).await?; |
| 410 | let from_key = repo.as_ref().map(store_key).unwrap_or_default(); |
| 411 | let (from_ns, _) = locate(&from_key); |
| 412 | let bound = self.store.namespaces(); |
| 413 | if let Some(why) = refusal(repo.as_ref(), &row.to_namespace, &bound, self.store.writable(&row.to_namespace), &from_ns) { |
| 414 | return Ok(Ran::Failed(why)); |
| 415 | } |
| 416 | let Some(repo) = repo else { return Ok(Ran::Failed("the repository is gone".to_owned())) }; |
| 417 | if !self.store.writable(&from_ns) { |
| 418 | return Ok(Ran::Later(format!("{from_ns} does not take writes now"))); |
| 419 | } |
| 420 | let forks = self.registry.all_forks_of(&repo.id).await?; |
| 421 | let mut items = vec![repo.clone()]; |
| 422 | items.extend(forks); |
| 423 | let ids: Vec<String> = items.iter().map(|item| item.id.clone()).collect(); |
| 424 | let reason = format!("moving to {}", row.to_namespace); |
| 425 | self.registry.pause(&ids, now_ms() + PAUSE_MS, &reason).await?; |
| 426 | match self.copy_and_switch(row, &items).await { |
| 427 | Ok(Ran::Moved) => Ok(Ran::Moved), |
| 428 | other => { |
| 429 | // Writes go on where they were. |
| 430 | self.registry.resume(&ids).await?; |
| 431 | other |
| 432 | } |
| 433 | } |
| 434 | } |
| 435 | |
| 436 | async fn copy_and_switch(&self, row: &MoveRow, items: &[Repo]) -> Result<Ran> { |
| 437 | // Push credentials already handed out reach the store directly. |
| 438 | let mut open_until = 0; |
| 439 | for item in items { |
| 440 | if let Some(read) = self.registry.by_id(&item.id).await? |
| 441 | && let Some(state) = crate::registry::refs_state(&read.id) |
| 442 | { |
| 443 | open_until = open_until.max(state.open_until); |
| 444 | } |
| 445 | } |
| 446 | let now = now_ms(); |
| 447 | if open_until > now { |
| 448 | if open_until - now > OPEN_WAIT_MS { |
| 449 | return Ok(Ran::Later("push credentials handed out have not expired yet".to_owned())); |
| 450 | } |
| 451 | worker::Delay::from(std::time::Duration::from_millis(open_until - now + 1_000)).await; |
| 452 | } |
| 453 | worker::Delay::from(std::time::Duration::from_millis(SETTLE_MS)).await; |
| 454 | |
| 455 | let default = self.store.default_namespace(); |
| 456 | let mut made: Vec<String> = Vec::new(); |
| 457 | let mut switched = Vec::new(); |
| 458 | let result = async { |
| 459 | for item in items { |
| 460 | let from = store_key(item); |
| 461 | let (_, name) = locate(&from); |
| 462 | let to = shards::compose(Some(&row.to_namespace), &name, &default); |
| 463 | // A removed working copy has nothing to copy: its row follows. |
| 464 | if crate::registry::retired(&item.id).is_some() { |
| 465 | switched.push((item.clone(), from, to, BTreeMap::new())); |
| 466 | continue; |
| 467 | } |
| 468 | self.store.create(&to, item.description.as_deref(), &item.default_branch).await?; |
| 469 | made.push(to.clone()); |
| 470 | match self.copy_until_same(&from, &to).await? { |
| 471 | Ok(refs) => switched.push((item.clone(), from, to, refs)), |
| 472 | Err(why) => return Ok(Err(why)), |
| 473 | } |
| 474 | } |
| 475 | Ok::<_, worker::Error>(Ok(())) |
| 476 | } |
| 477 | .await; |
| 478 | let failure = match result { |
| 479 | Ok(Ok(())) => None, |
| 480 | Ok(Err(why)) => Some(why), |
| 481 | Err(error) => Some(error.to_string()), |
| 482 | }; |
| 483 | if let Some(why) = failure { |
| 484 | // Nothing half made is left in the new namespace. |
| 485 | for key in made { |
| 486 | if let Err(error) = self.store.delete(&key).await { |
| 487 | worker::console_error!("move {}: {key} was left behind: {error}", row.id); |
| 488 | } |
| 489 | } |
| 490 | return Ok(Ran::Failed(why)); |
| 491 | } |
| 492 | self.registry.switch(&row.id, &switched, now_ms()).await?; |
| 493 | for (_, from, _, _) in &switched { |
| 494 | self.store.forget_access(from).await; |
| 495 | } |
| 496 | worker::console_log!("move {}: {} is in {} now", row.id, row.repo_id, row.to_namespace); |
| 497 | Ok(Ran::Moved) |
| 498 | } |
| 499 | |
| 500 | /// Copies `from`'s refs to `to` until both say the same and `from` has |
| 501 | /// not moved since; the refs as copied. |
| 502 | async fn copy_until_same(&self, from: &str, to: &str) -> Result<std::result::Result<BTreeMap<String, String>, String>> { |
| 503 | let source = self.store.open(from).await?.access(Scope::Read).await?; |
| 504 | let target_git = self.store.open(to).await?; |
| 505 | let target = target_git.access(Scope::Write).await?; |
| 506 | for _ in 0..ROUNDS { |
| 507 | let theirs = ref_map(refs::all(&source).await?); |
| 508 | let ours = ref_map(refs::all(&target).await?); |
| 509 | if theirs == ours { |
| 510 | return Ok(Ok(theirs)); |
| 511 | } |
| 512 | let (commands, wants, haves) = plan(&theirs, &ours); |
| 513 | if let Err(why) = land::copy_refs(&source, &target, &commands, &wants, &haves).await? { |
| 514 | return Ok(Err(format!("{to} refused the copy: {why}"))); |
| 515 | } |
| 516 | let again = ref_map(refs::all(&source).await?); |
| 517 | let copied = ref_map(refs::all(&target).await?); |
| 518 | if again == theirs && copied == theirs { |
| 519 | return Ok(Ok(theirs)); |
| 520 | } |
| 521 | } |
| 522 | Ok(Err(format!("{from} kept changing while it was copied"))) |
| 523 | } |
| 524 | |
| 525 | /// Deletes old copies kept past `MOVE_KEEP_DAYS`, each only while it |
| 526 | /// still says what was copied. |
| 527 | async fn clean_moves(&self) -> Result<()> { |
| 528 | let before = now_ms().saturating_sub(MOVE_KEEP_DAYS * DAY_MS); |
| 529 | for row in self.registry.moved_before(before, CLEANS_PER_SWEEP).await? { |
| 530 | let mut diverged = Vec::new(); |
| 531 | for copy in self.registry.copies_of(&row.id).await? { |
| 532 | let kept: BTreeMap<String, String> = serde_json::from_str(©.refs).unwrap_or_default(); |
| 533 | let now_refs = match self.store.open(©.from_key).await { |
| 534 | Ok(git) => match git.access(Scope::Read).await { |
| 535 | Ok(access) => refs::all(&access).await.map(ref_map), |
| 536 | Err(error) => Err(error), |
| 537 | }, |
| 538 | Err(error) => Err(error), |
| 539 | }; |
| 540 | let same = match now_refs { |
| 541 | Ok(refs) => refs == kept, |
| 542 | // Gone already. |
| 543 | Err(error) if error.to_string().contains("NOT_FOUND") => true, |
| 544 | Err(error) => return Err(error), |
| 545 | }; |
| 546 | if !same { |
| 547 | diverged.push(copy.from_key.clone()); |
| 548 | continue; |
| 549 | } |
| 550 | self.store.delete(©.from_key).await?; |
| 551 | self.registry.copy_cleaned(&row.id, ©.repo_id, now_ms()).await?; |
| 552 | } |
| 553 | if diverged.is_empty() { |
| 554 | self.registry.set_move(&row.id, "cleaned", None, now_ms()).await?; |
| 555 | } else { |
| 556 | let note = format!("kept, changed after the move: {}", diverged.join(", ")); |
| 557 | worker::console_error!("move {}: {note}", row.id); |
| 558 | self.registry.set_move(&row.id, "diverged", Some(¬e), now_ms()).await?; |
| 559 | } |
| 560 | } |
| 561 | Ok(()) |
| 562 | } |
| 563 | |
| 564 | /// `repo`, once writes to it no longer wait: the row is read again |
| 565 | /// every two seconds for up to [`PAUSE_WAIT_MS`]. `Err` with what to |
| 566 | /// tell the writer if they still wait. |
| 567 | pub(crate) async fn unpaused(&self, repo: Repo) -> Result<std::result::Result<Repo, (FailureCode, String)>> { |
| 568 | let started = now_ms(); |
| 569 | let mut repo = repo; |
| 570 | while let Some(reason) = crate::registry::paused(&repo.id, now_ms()) { |
| 571 | if now_ms().saturating_sub(started) >= PAUSE_WAIT_MS { |
| 572 | return Ok(Err((FailureCode::Conflict, paused_message(&repo, &reason)))); |
| 573 | } |
| 574 | worker::Delay::from(std::time::Duration::from_millis(PAUSE_POLL_MS)).await; |
| 575 | match self.registry.by_id(&repo.id).await? { |
| 576 | Some(read) => repo = read, |
| 577 | None => return Ok(Err((FailureCode::NotFound, "Repository not found.".to_owned()))), |
| 578 | } |
| 579 | } |
| 580 | Ok(Ok(repo)) |
| 581 | } |
| 582 | } |
| 583 | |
| 584 | #[cfg(test)] |
| 585 | mod tests { |
| 586 | use super::*; |
| 587 | |
| 588 | fn refs(pairs: &[(&str, &str)]) -> BTreeMap<String, String> { |
| 589 | pairs.iter().map(|(name, hash)| ((*name).to_owned(), (*hash).to_owned())).collect() |
| 590 | } |
| 591 | |
| 592 | #[test] |
| 593 | fn a_first_copy_sends_everything_and_a_second_only_what_changed() { |
| 594 | let a = "a".repeat(40); |
| 595 | let b = "b".repeat(40); |
| 596 | let c = "c".repeat(40); |
| 597 | let source = refs(&[("refs/heads/main", &a), ("refs/heads/dev", &b), ("refs/tags/v1", &a)]); |
| 598 | let (commands, wants, haves) = plan(&source, &BTreeMap::new()); |
| 599 | assert_eq!(commands.len(), 3); |
| 600 | assert!(commands.iter().all(|(_, old, _)| old.is_none())); |
| 601 | // Each object once. |
| 602 | assert_eq!(wants, vec![a.clone(), b.clone()]); |
| 603 | assert!(haves.is_empty()); |
| 604 | // A push landed during the copy: only it travels, with what the |
| 605 | // copy already holds named as had. |
| 606 | let moved = refs(&[("refs/heads/main", &c), ("refs/heads/dev", &b), ("refs/tags/v1", &a)]); |
| 607 | let (commands, wants, haves) = plan(&moved, &source); |
| 608 | assert_eq!(commands, vec![("refs/heads/main".to_owned(), Some(a.clone()), c.clone())]); |
| 609 | assert_eq!(wants, vec![c.clone()]); |
| 610 | assert_eq!(haves, vec![a.clone(), b.clone()]); |
| 611 | // A branch deleted during the copy is deleted from the copy too. |
| 612 | let deleted = refs(&[("refs/heads/main", &a)]); |
| 613 | let (commands, wants, _) = plan(&deleted, &source); |
| 614 | assert!(commands.iter().any(|(name, _, new)| name == "refs/heads/dev" && new == ZERO_ID)); |
| 615 | assert!(wants.is_empty()); |
| 616 | assert_eq!(plan(&source, &source).0, Vec::new()); |
| 617 | } |
| 618 | |
| 619 | #[test] |
| 620 | fn peeled_tags_are_not_refs() { |
| 621 | let map = ref_map(vec![("refs/tags/v1".into(), "t".into()), ("refs/tags/v1^{}".into(), "c".into()), ("refs/heads/main".into(), "m".into())]); |
| 622 | assert_eq!(map.len(), 2); |
| 623 | assert_eq!(map.get("refs/tags/v1").map(String::as_str), Some("t")); |
| 624 | } |
| 625 | |
| 626 | fn repo(fork_of: Option<&str>) -> Repo { |
| 627 | serde_json::from_value(serde_json::json!({ |
| 628 | "id": "rep_1", "namespace": "acme", "name": "rocket", "description": null, "isPrivate": true, |
| 629 | "ownerId": "usr_1", "defaultBranch": "main", "forkOf": fork_of, "protected": false, |
| 630 | "createdAt": "2026-10-07T00:00:00Z" |
| 631 | })) |
| 632 | .unwrap() |
| 633 | } |
| 634 | |
| 635 | #[test] |
| 636 | fn a_move_is_refused_when_it_cannot_work() { |
| 637 | let bound = vec!["g1t".to_owned(), "g1t-us-1".to_owned()]; |
| 638 | assert_eq!(refusal(Some(&repo(None)), "g1t-us-1", &bound, true, "g1t"), None); |
| 639 | assert!(refusal(None, "g1t-us-1", &bound, true, "g1t").unwrap().contains("gone")); |
| 640 | assert!(refusal(Some(&repo(Some("rep_0"))), "g1t-us-1", &bound, true, "g1t").unwrap().contains("working copy")); |
| 641 | assert!(refusal(Some(&repo(None)), "g1t-us-9", &bound, true, "g1t").unwrap().contains("not a bound")); |
| 642 | assert!(refusal(Some(&repo(None)), "g1t-us-1", &bound, false, "g1t").unwrap().contains("writes")); |
| 643 | assert!(refusal(Some(&repo(None)), "g1t", &bound, true, "g1t").unwrap().contains("already")); |
| 644 | assert!(paused_message(&repo(None), "moving to g1t-us-1").starts_with("acme/rocket is paused")); |
| 645 | } |
| 646 | |
| 647 | #[test] |
| 648 | fn writes_wait_less_than_a_move_may_take_and_the_pause_ends_on_its_own() { |
| 649 | assert!(PAUSE_WAIT_MS < PAUSE_MS); |
| 650 | // A handed-out credential lives five minutes plus a minute's grace. |
| 651 | assert!(OPEN_WAIT_MS > crate::store::CREDENTIAL_LIFE_MS + 60_000); |
| 652 | assert!(OPEN_WAIT_MS + SETTLE_MS < PAUSE_MS); |
| 653 | } |
| 654 | } |