Skip to content

g1t/services/repos/src/lib.rs

2,515 lines107,851 bytesCodeBlame
1//! The repos service: repository metadata, contents, forks, landing, and
2//! git over HTTPS.
3//!
4//! Other services reach it over `POST /rpc/<method>`; see
5//! `g1t_contracts::repos` for the methods and their arguments. Any other
6//! request is treated as git's smart HTTP protocol.
7
8mod backups;
9mod blame;
10mod catch_up;
11mod coalesce;
12mod commit_file;
13mod diff;
14mod fallback;
15mod forks;
16mod git_http;
17mod git_ops;
18mod import;
19mod land;
20mod last_commits;
21mod lifecycle;
22mod listing;
23mod meters;
24mod mirror;
25mod moves;
26mod namespaces;
27mod pack_cache;
28mod pack_limits;
29mod refs;
30mod refs_cache;
31mod registry;
32mod resilience;
33mod run_access;
34mod secret_scan;
35mod shards;
36mod shared;
37mod store;
38mod transfer;
39
40use g1t_contracts::events::{
41 Event, GitPush, NewEvent, Publish, RepoCreated, RepoForked, RepoUpdated, WorkspaceDeleted,
42 WorkspaceDeleting, WorkspaceRenamed, WorkspaceRestored,
43};
44use g1t_contracts::access::{self, Capability};
45use g1t_contracts::repos::*;
46use g1t_contracts::time::rfc3339;
47use g1t_contracts::{FailureCode, Outcome, PrincipalKind, User, Viewer, is_valid_repo_name, new_id};
48use g1t_kit::{args, now_ms, reply, rpc_method};
49use std::collections::{HashMap, HashSet, VecDeque};
50use std::rc::Rc;
51
52use serde::Serialize;
53use worker::{
54 Context, Env, Fetcher, MessageBatch, Method, Request, Response, Result, ScheduleContext, ScheduledEvent,
55 event,
56};
57
58use registry::{Registry, can_read, can_write, store_key};
59use store::{ArtifactsStore, GitRepo, GitStore, Scope};
60
61/// Namespace that holds every pull request's fork: `pulls/<pull id>`.
62pub(crate) const PULLS_NAMESPACE: &str = "pulls";
63const MAX_TEXT_BYTES: usize = 512 * 1024;
64/// How far back a pull request may have forked and still be landed.
65const MAX_ANCESTRY: u32 = 1000;
66/// The most tags a repository's Tags page reads and lists.
67const MAX_TAGS_READ: usize = 100;
68
69/// One path segment, percent-encoded for a cache key.
70fn urlencoding_segment(segment: &str) -> String {
71 segment
72 .bytes()
73 .map(|b| if b.is_ascii_alphanumeric() || b"-._~".contains(&b) { (b as char).to_string() } else { format!("%{b:02X}") })
74 .collect()
75}
76const MAX_DESCRIPTION_CHARS: usize = 200;
77pub(crate) const SOURCE: &str = "repos";
78pub(crate) const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
79
80pub(crate) fn not_found<T>() -> Outcome<T> {
81 Outcome::fail(FailureCode::NotFound, "Repository not found.")
82}
83
84/// Decoded text, or `None` when the file is too large or looks binary.
85/// Whether a ref is a full commit hash rather than a branch name.
86fn is_commit_hash(git_ref: &str) -> bool {
87 git_ref.len() == 40 && git_ref.bytes().all(|b| b.is_ascii_hexdigit())
88}
89
90fn text_of(bytes: Vec<u8>) -> Option<String> {
91 if bytes.len() > MAX_TEXT_BYTES || bytes.contains(&0) {
92 return None;
93 }
94 Some(String::from_utf8_lossy(&bytes).into_owned())
95}
96
97fn is_readme(name: &str) -> bool {
98 matches!(
99 name.to_lowercase().as_str(),
100 "readme" | "readme.md" | "readme.markdown" | "readme.txt"
101 )
102}
103
104/// Whether `ancestor` is reachable from the newest commit in `history`.
105///
106/// `history` is the first-parent chain, which is all the store lists; a fork
107/// that merged the target branch in has the target's head on a second
108/// parent, so the walk follows every parent.
109async fn descends_from<R: GitRepo>(repo: &R, history: &[Commit], ancestor: &str) -> Result<bool> {
110 let known: HashMap<&str, &[String]> = history
111 .iter()
112 .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
113 .collect();
114 let mut seen = HashSet::new();
115 let mut queue: Vec<String> = history
116 .first()
117 .map(|c| c.hash.clone())
118 .into_iter()
119 .collect();
120 while let Some(hash) = queue.pop() {
121 if hash == ancestor {
122 return Ok(true);
123 }
124 if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
125 continue;
126 }
127 match known.get(hash.as_str()) {
128 Some(parents) => queue.extend(parents.iter().cloned()),
129 None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
130 }
131 }
132 Ok(false)
133}
134
135/// The commit closest to the newest in `history` that is also in `shared`:
136/// where a fork and the repository it came from last agreed.
137async fn nearest_ancestor_in<R: GitRepo>(
138 repo: &R,
139 history: &[Commit],
140 shared: &HashSet<String>,
141) -> Result<Option<String>> {
142 let known: HashMap<&str, &[String]> = history
143 .iter()
144 .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
145 .collect();
146 let mut seen = HashSet::new();
147 let mut queue: VecDeque<String> = history
148 .first()
149 .map(|c| c.hash.clone())
150 .into_iter()
151 .collect();
152 while let Some(hash) = queue.pop_front() {
153 if shared.contains(&hash) {
154 return Ok(Some(hash));
155 }
156 if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
157 continue;
158 }
159 match known.get(hash.as_str()) {
160 Some(parents) => queue.extend(parents.iter().cloned()),
161 None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
162 }
163 }
164 Ok(None)
165}
166
167thread_local! {
168 /// Targets' sides of mergeability, by head (coalesce.rs).
169 static TARGETS: std::cell::RefCell<coalesce::Memo<coalesce::TargetKey, Rc<coalesce::TargetSide>>> =
170 std::cell::RefCell::new(coalesce::Memo::new(coalesce::TARGET_TTL_MS, 32));
171 /// What targets changed between two trees.
172 static THEIRS: std::cell::RefCell<coalesce::Memo<coalesce::TheirsKey, (Vec<String>, bool)>> =
173 std::cell::RefCell::new(coalesce::Memo::new(coalesce::THEIRS_TTL_MS, 256));
174 /// What repositories hold, as read for a push's first request, for the
175 /// same push's second: a push's POST does not wait on the database.
176 static HELD: std::cell::RefCell<coalesce::Memo<String, u64>> =
177 std::cell::RefCell::new(coalesce::Memo::new(60_000, 512));
178}
179
180pub(crate) struct Repos<S: GitStore> {
181 registry: Registry,
182 store: S,
183 events: Fetcher,
184 /// Asked during a push which secrets have been allowed.
185 security: Option<Fetcher>,
186 /// Asked whether a workspace is on a plan, for its private storage.
187 billing: Option<Fetcher>,
188 /// Told when a repository moves, for the tokens of agents at work on it.
189 identity: Option<Fetcher>,
190 /// What a free workspace's private repositories may hold.
191 free_private_bytes: i64,
192 /// Days a pull request's working copy is kept after it settles (forks.rs).
193 pub(crate) fork_days: u64,
194 /// The most a repository may hold (pack_limits.rs), and what happens
195 /// to a push too large to scan.
196 repo_limit: u64,
197 large_pushes: git_http::LargePushes,
198 /// Which git store namespace new repositories go in (shards.rs), and
199 /// the most each should hold (`ARTIFACTS_NAMESPACE_LIMITS`).
200 placement: shards::Placement,
201 limits: HashMap<String, u64>,
202 /// What isolates share: answers that list refs (refs_cache.rs).
203 shared: Option<Rc<shared::Shared>>,
204 /// Packs for fresh clones (pack_cache.rs); `None` without the bucket.
205 packs: Option<Rc<pack_cache::Packs>>,
206}
207
208impl<S: GitStore> Repos<S> {
209 /// Records that the refs of the repository with this id changed, once
210 /// they have, so that the answers kept that list them go stale (see
211 /// refs_cache.rs). Everything that changes a repository's refs calls
212 /// this after it (`every_ref_writer_records_the_change` checks). A
213 /// failure is logged: the change itself happened, and what was kept
214 /// expires within `refs_cache::TTL_SECONDS` regardless.
215 pub(crate) async fn refs_moved(&self, repo_id: &str) {
216 if let Err(error) = self.registry.refs_moved(repo_id).await {
217 worker::console_error!("refs of {repo_id} changed but not recorded: {error}");
218 }
219 }
220
221 pub(crate) async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> {
222 g1t_kit::call(
223 &self.events,
224 "publish",
225 &Publish {
226 events: vec![event],
227 },
228 )
229 .await
230 }
231
232 /// Whether the viewer may read `repo`. A pull request's fork of a
233 /// private repository can be read by everyone who can read that
234 /// repository, so its members can review and check out the change, as
235 /// well as by whoever opened the pull request.
236 async fn may_read(&self, repo: &Repo, viewer: &Viewer) -> Result<bool> {
237 if can_read(repo, viewer) {
238 return Ok(true);
239 }
240 let Some(source_id) = &repo.fork_of else {
241 return Ok(false);
242 };
243 Ok(self
244 .registry
245 .by_id(source_id)
246 .await?
247 .is_some_and(|source| can_read(&source, viewer)))
248 }
249
250 /// `repo`, if there is one and the viewer may read it.
251 async fn visible(&self, repo: Option<Repo>, viewer: &Viewer) -> Result<Option<Repo>> {
252 Ok(match repo {
253 Some(repo) if self.may_read(&repo, viewer).await? => Some(repo),
254 _ => None,
255 })
256 }
257
258 /// Resolves a repo the viewer may read; private repos look missing.
259 pub(crate) async fn readable(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
260 self.visible(self.registry.by_path(path).await?, viewer)
261 .await
262 }
263
264 async fn get(&self, a: GetArgs) -> Result<Outcome<Repo>> {
265 Ok(self
266 .readable(&a.path, &a.viewer)
267 .await?
268 .map_or_else(not_found, Outcome::Ok))
269 }
270
271 async fn get_by_id(&self, a: GetByIdArgs) -> Result<Outcome<Repo>> {
272 Ok(self
273 .visible(self.registry.by_id(&a.id).await?, &a.viewer)
274 .await?
275 .map_or_else(not_found, Outcome::Ok))
276 }
277
278 async fn update(&self, a: UpdateArgs) -> Result<Outcome<Repo>> {
279 let viewer = Some(a.actor.clone());
280 let Some(repo) = self.readable(&a.path, &viewer).await? else {
281 return Ok(not_found());
282 };
283 // Its details take Maintain; its protection, Maintain too; who can
284 // see it, Admin (below). See g1t_contracts::access.
285 let protection_changes = a.protected.is_some_and(|protected| protected != repo.protected);
286 let details_change = a.description.is_some() || a.website.is_some() || a.topics.is_some();
287 let mut needed = Vec::new();
288 if details_change || !protection_changes {
289 needed.push(Capability::ManageSettings);
290 }
291 if protection_changes {
292 needed.push(Capability::ManageProtection);
293 }
294 let full_name = format!("{}/{}", repo.namespace, repo.name);
295 if repo.fork_of.is_some() {
296 return Ok(Outcome::fail(FailureCode::Forbidden, access::needs(Capability::ManageSettings, &full_name)));
297 }
298 if let Some(missing) = needed.into_iter().find(|capability| !registry::can(&repo, &viewer, *capability)) {
299 return Ok(Outcome::fail(FailureCode::Forbidden, access::needs(missing, &full_name)));
300 }
301 if !a.actor.verified {
302 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
303 }
304 if let Some((code, message)) = lifecycle::archived_refusal(&repo) {
305 return Ok(Outcome::fail(code, message));
306 }
307 let description = match a.description {
308 Some(text) => Some(
309 text.trim()
310 .chars()
311 .take(MAX_DESCRIPTION_CHARS)
312 .collect::<String>(),
313 )
314 .filter(|text| !text.is_empty()),
315 None => repo.description.clone(),
316 };
317 let website = match a.website.as_deref() {
318 Some(text) => match clean_website(text) {
319 Ok(website) => website,
320 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
321 },
322 None => repo.website.clone(),
323 };
324 // Who can see it is an owner's to change, and a free workspace's
325 // storage may not take it private: see lifecycle.rs.
326 let wants_private = a.is_private.filter(|private| *private != repo.is_private);
327 if wants_private.is_some()
328 && let Err((code, message)) = lifecycle::admin_only(
329 lifecycle::Asker::on(&a.actor, &repo),
330 &repo.namespace,
331 "change the visibility of",
332 Capability::Administer,
333 )
334 {
335 return Ok(Outcome::fail(code, message));
336 }
337 let is_private = repo.is_private;
338 let protected = a.protected.unwrap_or(repo.protected);
339 let topics = match &a.topics {
340 Some(topics) => match clean_topics(topics) {
341 Ok(topics) => topics,
342 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
343 },
344 None => repo.topics.clone(),
345 };
346 self.registry
347 .update(&repo.id, description.as_deref(), protected, &topics, website.as_deref())
348 .await?;
349 let updated = Repo {
350 description,
351 is_private,
352 protected,
353 topics,
354 website,
355 ..repo
356 };
357 if let Some(private) = wants_private {
358 return self.change_visibility(updated, private, &a.actor, a.surface).await;
359 }
360 let visibility_changed = false;
361 // Search and anything else that shows the repository hears of it;
362 // a change of visibility is announced on its own as well, so that
363 // what was public stops being shown at once.
364 self.publish(NewEvent {
365 kind: "repo.updated",
366 source: SOURCE,
367 repo_id: Some(updated.id.clone()),
368 actor: Some(a.actor.id.clone()),
369 data: RepoUpdated {
370 repo_id: updated.id.clone(),
371 namespace: updated.namespace.clone(),
372 name: updated.name.clone(),
373 is_private,
374 visibility_changed,
375 },
376 })
377 .await?;
378 Ok(Outcome::Ok(updated))
379 }
380
381 /// The repository with this id, if it is not a fork, and its store.
382 async fn stored(&self, repo_id: &str) -> Result<Option<S::Repo>> {
383 match self.registry.by_id(repo_id).await? {
384 Some(repo) if repo.fork_of.is_none() => Ok(Some(self.store.open(&store_key(&repo)).await?)),
385 _ => Ok(None),
386 }
387 }
388
389 async fn list_files(&self, a: ListFilesArgs) -> Result<FileList> {
390 let Some(repo) = self.registry.by_id(&a.repo_id).await?.filter(|repo| repo.fork_of.is_none()) else {
391 return Ok(FileList::default());
392 };
393 let git = self.store.open(&store_key(&repo)).await?;
394 let head = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
395 listing::list(&git, None, &head, &a.skip_dirs, a.limit).await
396 }
397
398 async fn changed_files(&self, a: ChangedFilesArgs) -> Result<FileList> {
399 let Some(git) = self.stored(&a.repo_id).await? else {
400 return Ok(FileList::default());
401 };
402 listing::list(&git, a.base.as_deref(), &a.head, &a.skip_dirs, a.limit).await
403 }
404
405 /// Branches and tags with their commits, for g1t's own services.
406 async fn refs_of(&self, a: RefsArgs) -> Result<Option<RepoRefs>> {
407 let Some(repo) = self.registry.by_id(&a.repo_id).await?.filter(|repo| repo.fork_of.is_none()) else {
408 return Ok(None);
409 };
410 let git = self.store.open(&store_key(&repo)).await?;
411 let access = git.access(Scope::Read).await?;
412 let refs = refs::heads_and_tags(refs::all(&access).await?)
413 .into_iter()
414 .map(|(name, commit)| GitRefEntry { name, commit })
415 .collect();
416 Ok(Some(RepoRefs { repo, refs }))
417 }
418
419 async fn raw_file(&self, a: RawFileArgs) -> Result<Option<RawFile>> {
420 use base64::Engine;
421 let Some(git) = self.stored(&a.repo_id).await? else {
422 return Ok(None);
423 };
424 Ok(git
425 .read_file(&a.git_ref, &a.path)
426 .await?
427 .filter(|bytes| bytes.len() <= a.max_bytes as usize)
428 .map(|bytes| RawFile { size: bytes.len() as u64, data: base64::engine::general_purpose::STANDARD.encode(bytes) }))
429 }
430
431 async fn raw_blobs(&self, a: RawBlobsArgs) -> Result<Vec<RawBlob>> {
432 use base64::Engine;
433 let Some(git) = self.stored(&a.repo_id).await? else {
434 return Ok(Vec::new());
435 };
436 let hashes: Vec<&String> = a.hashes.iter().take(MAX_READ_BLOBS).collect();
437 let mut out = Vec::with_capacity(hashes.len());
438 // A few at a time, as listing::read does: each is a round trip.
439 for group in hashes.chunks(8) {
440 let read = futures_util::future::try_join_all(group.iter().map(|hash| git.read_blob(hash))).await?;
441 for (hash, bytes) in group.iter().zip(read) {
442 let size = bytes.as_ref().map_or(0, |bytes| bytes.len() as u64);
443 let data = bytes
444 .filter(|bytes| bytes.len() <= a.max_bytes as usize)
445 .map(|bytes| base64::engine::general_purpose::STANDARD.encode(bytes));
446 out.push(RawBlob { hash: (*hash).clone(), size, data });
447 }
448 }
449 Ok(out)
450 }
451
452 async fn read_blobs(&self, a: ReadBlobsArgs) -> Result<Vec<BlobText>> {
453 let Some(git) = self.stored(&a.repo_id).await? else {
454 return Ok(Vec::new());
455 };
456 listing::read(&git, &a.hashes, a.max_bytes.min(MAX_TEXT_BYTES as u32)).await
457 }
458
459 async fn create(&self, a: CreateArgs) -> Result<Outcome<Repo>> {
460 if !a.owner.verified {
461 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
462 }
463 let name = a.name.trim().to_lowercase();
464 if !is_valid_repo_name(&name) {
465 return Ok(Outcome::fail(
466 FailureCode::Invalid,
467 "Use letters, digits, dots, hyphens and underscores only.",
468 ));
469 }
470 let namespace = a.namespace.trim().to_lowercase();
471 if namespace.is_empty() {
472 return Ok(Outcome::fail(
473 FailureCode::Invalid,
474 "Say which workspace to create the repository in.",
475 ));
476 }
477 if !a.owner.is_member(&namespace) {
478 return Ok(Outcome::fail(
479 FailureCode::Forbidden,
480 "You are not a member of that workspace.",
481 ));
482 }
483 let path = RepoPath { namespace, name };
484 match self.registry.by_path_any(&path).await? {
485 Some((_, None)) => {
486 return Ok(Outcome::fail(
487 FailureCode::Conflict,
488 "That workspace already has a repository with that name.",
489 ));
490 }
491 Some((_, Some(_))) => {
492 return Ok(Outcome::fail(
493 FailureCode::Conflict,
494 format!(
495 "{}/{} was deleted recently and can still be restored, so its name is taken. Restore it, or delete it permanently from the workspace's Recently deleted list.",
496 path.namespace, path.name
497 ),
498 ));
499 }
500 None => {}
501 }
502 // With a credential (a GitHub App installation's token), everything
503 // is copied: every branch and tag. See mirror.rs.
504 let mut credentialed = None;
505 if let (Some(url), Some(token)) = (a.import_url.as_deref(), a.import_token.as_deref()) {
506 let Some(url) = import::clean_url(url) else {
507 return Ok(Outcome::fail(FailureCode::Invalid, "That is not an https repository address."));
508 };
509 let source = mirror::Endpoint::github(&url, token);
510 match mirror::probe(&source).await? {
511 Ok(advertised) => credentialed = Some((source, advertised)),
512 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
513 }
514 }
515 // An import is fetched before anything is created, so that an
516 // address that does not work leaves nothing behind.
517 let mut imported = None;
518 if let Some(url) = a
519 .import_url
520 .as_deref()
521 .map(str::trim)
522 .filter(|url| !url.is_empty() && credentialed.is_none())
523 {
524 let Some(url) = import::clean_url(url) else {
525 return Ok(Outcome::fail(
526 FailureCode::Invalid,
527 "Give the https address of a public repository, such as https://github.com/owner/repo.",
528 ));
529 };
530 let remote = match import::discover(&url).await? {
531 Ok(remote) => remote,
532 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
533 };
534 imported = Some((remote, url));
535 }
536 let now = now_ms();
537 let repo = Repo {
538 id: new_id("rep", now),
539 namespace: path.namespace,
540 name: path.name,
541 description: a
542 .description
543 .map(|text| text.trim().to_owned())
544 .filter(|text| !text.is_empty()),
545 is_private: a.is_private,
546 owner_id: a.owner.id.clone(),
547 default_branch: imported
548 .as_ref()
549 .map(|(remote, _)| remote.branch.clone())
550 .or_else(|| credentialed.as_ref().and_then(|(_, advertised)| advertised.default_branch()))
551 .unwrap_or_else(|| "main".to_owned()),
552 fork_of: None,
553 protected: false,
554 created_at: rfc3339(now),
555 topics: Vec::new(),
556 website: None,
557 archived_at: None,
558 };
559 let namespace = match self.place(&repo).await? {
560 Ok(namespace) => namespace,
561 Err(unplaced) => return Ok(Outcome::fail(FailureCode::Conflict, unplaced.message())),
562 };
563 self.registry
564 .claim_store_key(&repo, namespace.as_deref(), &self.store.default_namespace())
565 .await?;
566 self.store
567 .create(
568 &store_key(&repo),
569 repo.description.as_deref(),
570 &repo.default_branch,
571 )
572 .await?;
573 self.registry.insert(&repo).await?;
574 // A repository that was transferred away from this path stops
575 // redirecting here.
576 self.registry
577 .drop_redirect(&RepoPath {
578 namespace: repo.namespace.clone(),
579 name: repo.name.clone(),
580 })
581 .await?;
582 // Every branch and tag the import made, announced as pushes.
583 let mut pushed: Vec<(String, String)> = Vec::new();
584 // A public repository, read with no credential: every branch and
585 // tag is copied too, the default branch the one its HEAD names.
586 if let Some((_, url)) = imported {
587 let access = self
588 .store
589 .open(&store_key(&repo))
590 .await?
591 .access(Scope::Write)
592 .await?;
593 let target = mirror::Endpoint::bearer(&access.remote, &access.token);
594 let copied = mirror::copy(&mirror::Endpoint::anonymous(&url), &target, mirror::Prune::Yes).await?;
595 self.refs_moved(&repo.id).await;
596 match copied {
597 Ok(copied) => pushed = mirror::import_pushes(&copied.updated, &repo.default_branch),
598 Err(reason) => {
599 self.registry.remove(&repo.id).await?;
600 return Ok(Outcome::fail(
601 FailureCode::Invalid,
602 format!("The repository could not be stored: {reason}"),
603 ));
604 }
605 }
606 }
607 if let Some((source, _)) = credentialed {
608 let access = self
609 .store
610 .open(&store_key(&repo))
611 .await?
612 .access(Scope::Write)
613 .await?;
614 let target = mirror::Endpoint::bearer(&access.remote, &access.token);
615 let copied = mirror::copy(&source, &target, mirror::Prune::Yes).await?;
616 self.refs_moved(&repo.id).await;
617 match copied {
618 Ok(copied) => pushed = mirror::import_pushes(&copied.updated, &repo.default_branch),
619 Err(reason) => {
620 self.registry.remove(&repo.id).await?;
621 return Ok(Outcome::fail(
622 FailureCode::Invalid,
623 format!("The repository could not be copied: {reason}"),
624 ));
625 }
626 }
627 }
628 self.publish(NewEvent {
629 kind: "repo.created",
630 source: SOURCE,
631 repo_id: Some(repo.id.clone()),
632 actor: Some(a.owner.id),
633 data: RepoCreated {
634 repo_id: repo.id.clone(),
635 namespace: repo.namespace.clone(),
636 name: repo.name.clone(),
637 is_private: repo.is_private,
638 },
639 })
640 .await?;
641 for (git_ref, head) in &pushed {
642 self.publish_push(&repo, git_ref, None, head, None).await?;
643 }
644 Ok(Outcome::Ok(repo))
645 }
646
647 /// Where a workspace keeps its data, asked of identity only when an EU
648 /// namespace is configured: without one, every workspace's
649 /// repositories go anywhere and identity is never asked.
650 async fn residency_of(&self, workspace: &str) -> Result<shards::Residency> {
651 if self.placement.eu.is_none() {
652 return Ok(shards::Residency::Anywhere);
653 }
654 let Some(identity) = &self.identity else {
655 return Ok(shards::Residency::Anywhere);
656 };
657 let residency: Option<g1t_contracts::identity::DataResidency> = g1t_kit::call(
658 identity,
659 "workspace_residency",
660 &g1t_contracts::identity::SlugArgs { slug: workspace.to_owned() },
661 )
662 .await?;
663 Ok(match residency {
664 Some(g1t_contracts::identity::DataResidency::Eu) => shards::Residency::Eu,
665 _ => shards::Residency::Anywhere,
666 })
667 }
668
669 /// How each bound namespace stands (namespaces.rs).
670 async fn standings(&self) -> Result<Vec<namespaces::Standing>> {
671 let bound = self.store.namespaces();
672 let default = self.store.default_namespace();
673 let now = now_ms();
674 let config = namespaces::Configured {
675 bound: &bound,
676 default: &default,
677 placement: &self.placement,
678 limits: &self.limits,
679 on_fallback: &|namespace| self.store.on_fallback(&shards::compose(Some(namespace), "x", &default)),
680 writable: &|namespace| self.store.writable(namespace),
681 breaker_open: &|namespace| resilience::open_now(namespace, now),
682 };
683 let (held, recent) = futures_util::future::join(namespaces::held(&self.registry.db), namespaces::recent(&self.registry.db, now)).await;
684 Ok(namespaces::standings(&config, &held?, &recent?))
685 }
686
687 /// The namespace a new repository goes in (shards.rs): its workspace's
688 /// residency, then how each namespace stands, read only when there is
689 /// a choice to make. `Ok(None)` for the default.
690 async fn place(&self, repo: &Repo) -> Result<std::result::Result<Option<String>, shards::Unplaced>> {
691 let residency = self.residency_of(&repo.namespace).await?;
692 let bound = self.store.namespaces();
693 let loads = if residency == shards::Residency::Anywhere && !self.placement.needs_loads(&bound) {
694 // One namespace to choose from at most: nothing to read.
695 bound
696 .iter()
697 .map(|namespace| shards::Load {
698 namespace: namespace.clone(),
699 bound: true,
700 writable: self.store.writable(namespace),
701 ..shards::Load::default()
702 })
703 .collect()
704 } else {
705 let default = self.store.default_namespace();
706 let now = now_ms();
707 let config = namespaces::Configured {
708 bound: &bound,
709 default: &default,
710 placement: &self.placement,
711 limits: &self.limits,
712 on_fallback: &|namespace| self.store.on_fallback(&shards::compose(Some(namespace), "x", &default)),
713 writable: &|namespace| self.store.writable(namespace),
714 breaker_open: &|namespace| resilience::open_now(namespace, now),
715 };
716 namespaces::loads(&self.registry.db, &config, now).await?
717 };
718 Ok(self.placement.choose(&repo.id, residency, &loads))
719 }
720
721 /// `storage_options`: what a workspace may choose about where its
722 /// repositories are kept.
723 fn storage_options(&self) -> StorageOptions {
724 let bound = self.store.namespaces();
725 StorageOptions { eu_available: self.placement.eu_available(&bound, |namespace| self.store.writable(namespace)) }
726 }
727
728 async fn tree(&self, a: TreeArgs) -> Result<Outcome<TreeView>> {
729 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
730 return Ok(not_found());
731 };
732 let git = self.read_git(&repo).await?;
733 let git_ref = a
734 .git_ref
735 .clone()
736 .unwrap_or_else(|| repo.default_branch.clone());
737
738 let Some(head) = git.log(&git_ref, 1).await?.into_iter().next() else {
739 // An unknown ref is an error; a repo with no commits is just empty.
740 if a.git_ref.is_some() {
741 return Ok(Outcome::fail(
742 FailureCode::NotFound,
743 "No such branch, tag or commit.",
744 ));
745 }
746 return Ok(Outcome::Ok(TreeView {
747 repo,
748 git_ref,
749 path: a.tree_path,
750 head: None,
751 entries: Vec::new(),
752 readme: None,
753 }));
754 };
755
756 let no_directory = || Outcome::fail(FailureCode::NotFound, "No such directory.");
757 let mut entries = git.read_tree(&head.tree_hash).await?;
758 for segment in a.tree_path.split('/').filter(|segment| !segment.is_empty()) {
759 let next = entries.as_ref().and_then(|entries| {
760 entries
761 .iter()
762 .find(|entry| entry.name == segment && entry.kind == EntryKind::Tree)
763 });
764 let Some(next) = next else {
765 return Ok(no_directory());
766 };
767 entries = git.read_tree(&next.hash).await?;
768 }
769 let Some(mut entries) = entries else {
770 return Ok(no_directory());
771 };
772 // Directories first, then by name.
773 entries.sort_by(|a, b| {
774 (b.kind == EntryKind::Tree)
775 .cmp(&(a.kind == EntryKind::Tree))
776 .then_with(|| a.name.cmp(&b.name))
777 });
778
779 let readme_entry = entries
780 .iter()
781 .find(|entry| entry.kind == EntryKind::Blob && is_readme(&entry.name));
782 let readme = match readme_entry {
783 Some(entry) => git.read_blob(&entry.hash).await?.map(|bytes| Readme {
784 name: entry.name.clone(),
785 text: text_of(bytes),
786 }),
787 None => None,
788 };
789 Ok(Outcome::Ok(TreeView {
790 repo,
791 git_ref,
792 path: a.tree_path,
793 head: Some(head),
794 entries,
795 readme,
796 }))
797 }
798
799 async fn blob(&self, a: BlobArgs) -> Result<Outcome<BlobView>> {
800 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
801 return Ok(not_found());
802 };
803 let bytes = if a.file_path.is_empty() {
804 None
805 } else {
806 let git = self.read_git(&repo).await?;
807 git.read_file(&a.git_ref, &a.file_path).await?
808 };
809 let Some(bytes) = bytes else {
810 return Ok(Outcome::fail(FailureCode::NotFound, "No such file."));
811 };
812 Ok(Outcome::Ok(BlobView {
813 repo,
814 git_ref: a.git_ref,
815 path: a.file_path,
816 size: bytes.len() as u64,
817 text: text_of(bytes),
818 }))
819 }
820
821 async fn blame(&self, a: BlameArgs) -> Result<Outcome<Blame>> {
822 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
823 return Ok(not_found());
824 };
825 let git = self.read_git(&repo).await?;
826 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
827 Ok(match blame::blame(&git, &git_ref, &a.file_path).await? {
828 Some(blame) => Outcome::Ok(blame),
829 None => not_found(),
830 })
831 }
832
833 async fn log(&self, a: LogArgs) -> Result<Outcome<Vec<Commit>>> {
834 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
835 return Ok(not_found());
836 };
837 let git = self.read_git(&repo).await?;
838 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
839 Ok(Outcome::Ok(git.log(&git_ref, a.limit).await?))
840 }
841
842 /// Which commit last changed each entry of a directory. Kept in this
843 /// colo's cache by repository, head commit and path: a commit's history
844 /// never changes, so an answer is good for as long as it is kept.
845 async fn last_commits(&self, a: g1t_contracts::repos::LastCommitsArgs) -> Result<Outcome<g1t_contracts::repos::LastCommits>> {
846 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
847 return Ok(not_found());
848 };
849 let git = self.read_git(&repo).await?;
850 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
851 let Some(head) = git.log(&git_ref, 1).await?.into_iter().next() else {
852 return Ok(Outcome::fail(FailureCode::NotFound, "No such branch, tag or commit."));
853 };
854 let key = format!(
855 "https://last-commits.g1t.internal/{}/{}/{}",
856 repo.id,
857 head.hash,
858 a.tree_path.split('/').map(urlencoding_segment).collect::<Vec<_>>().join("/")
859 );
860 let cache = worker::Cache::default();
861 if let Ok(Some(mut kept)) = cache.get(key.as_str(), false).await {
862 if let Ok(found) = kept.json::<g1t_contracts::repos::LastCommits>().await {
863 return Ok(Outcome::Ok(found));
864 }
865 }
866 // Asked with a budget: past it, what was found so far, not kept.
867 let started = worker::Date::now().as_millis();
868 let budget = a.budget_ms;
869 let out_of_time = move || budget.is_some_and(|budget| worker::Date::now().as_millis().saturating_sub(started) > budget);
870 let (entries, complete) = last_commits::last_commits(&git, &head.hash, &a.tree_path, &out_of_time).await?;
871 let stopped = out_of_time();
872 let found = g1t_contracts::repos::LastCommits { entries, complete };
873 if stopped && !found.complete {
874 return Ok(Outcome::Ok(found));
875 }
876 if let Ok(mut response) = worker::Response::from_json(&found) {
877 let _ = response.headers_mut().set("cache-control", "max-age=604800");
878 let _ = cache.put(key.as_str(), response).await;
879 }
880 Ok(Outcome::Ok(found))
881 }
882
883 /// The repository's tags, newest commit first, at most 100.
884 async fn tags(&self, a: g1t_contracts::repos::TagsArgs) -> Result<Outcome<Vec<g1t_contracts::repos::Tag>>> {
885 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
886 return Ok(not_found());
887 };
888 let git = self.store.open(&store_key(&repo)).await?;
889 let access = git.access(Scope::Read).await?;
890 let named: Vec<(String, String)> = refs::heads_and_tags(refs::all(&access).await?)
891 .into_iter()
892 .filter_map(|(name, hash)| name.strip_prefix("refs/tags/").map(|tag| (tag.to_owned(), hash)))
893 .collect();
894 let read = self.read_git(&repo).await?;
895 let commits = futures_util::future::join_all(named.iter().take(MAX_TAGS_READ).map(|(_, hash)| read.log(hash, 1))).await;
896 let mut tags: Vec<g1t_contracts::repos::Tag> = named
897 .into_iter()
898 .zip(commits.into_iter().map(|found| found.ok().and_then(|list| list.into_iter().next())).chain(std::iter::repeat(None)))
899 .map(|((name, _), commit)| g1t_contracts::repos::Tag { name, commit })
900 .collect();
901 tags.sort_by(|a, b| {
902 let at = |tag: &g1t_contracts::repos::Tag| tag.commit.as_ref().map(|c| c.authored_at.clone()).unwrap_or_default();
903 at(b).cmp(&at(a)).then_with(|| b.name.cmp(&a.name))
904 });
905 tags.truncate(MAX_TAGS_READ);
906 Ok(Outcome::Ok(tags))
907 }
908
909 /// The repository's branches, default branch first.
910 async fn branches(&self, a: BranchesArgs) -> Result<Outcome<Vec<Branch>>> {
911 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
912 return Ok(not_found());
913 };
914 let mut branches = self.read_git(&repo).await?.branches().await?;
915 branches.sort_by_key(|branch| branch.name != repo.default_branch);
916 Ok(Outcome::Ok(branches))
917 }
918
919 /// Whether a pull request's source lacks commits that the branch it
920 /// would merge into has.
921 async fn behind(&self, a: BehindArgs) -> Result<bool> {
922 let Some(source) = self.registry.by_id(&a.source_id).await? else {
923 return Ok(false);
924 };
925 let target = match &source.fork_of {
926 Some(id) => self.registry.by_id(id).await?,
927 None => Some(source.clone()),
928 };
929 let Some(target) = target else {
930 return Ok(false);
931 };
932 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
933 let target_branch = a.target_branch.unwrap_or_else(|| target.default_branch.clone());
934 let target_head = self
935 .read_git(&target)
936 .await?
937 .log(&target_branch, 1)
938 .await?
939 .into_iter()
940 .next()
941 .map(|commit| commit.hash);
942 let Some(target_head) = target_head else {
943 return Ok(false);
944 };
945 let source_git = self.read_git(&source).await?;
946 let history = source_git.log(&branch, MAX_ANCESTRY).await?;
947 if history.is_empty() {
948 return Ok(false);
949 }
950 Ok(!descends_from(&source_git, &history, &target_head).await?)
951 }
952
953 /// The files a pull request's source and the default branch it would
954 /// merge into each changed since they last agreed. Where the two lists
955 /// share no file, the merge cannot conflict; where they do, it may.
956 async fn divergence(&self, a: BehindArgs) -> Result<Option<Divergence>> {
957 let Some(source) = self.registry.by_id(&a.source_id).await? else {
958 return Ok(None);
959 };
960 let target = match &source.fork_of {
961 Some(id) => self.registry.by_id(id).await?,
962 None => Some(source.clone()),
963 };
964 let Some(target) = target else {
965 return Ok(None);
966 };
967 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
968 let target_branch = a.target_branch.unwrap_or_else(|| target.default_branch.clone());
969 let source_git = self.read_git(&source).await?;
970 let target_git = self.read_git(&target).await?;
971 // The target's side is the same for every pull request into it, and
972 // worked out once per head (coalesce.rs).
973 let (history, side) = futures_util::future::try_join(
974 source_git.log(&branch, MAX_ANCESTRY),
975 self.target_side(&target, &target_git, &target_branch),
976 )
977 .await?;
978 let target_history = &side.history;
979 let (Some(head), Some(base)) = (history.first(), target_history.first()) else {
980 return Ok(None);
981 };
982 let behind = !descends_from(&source_git, &history, &base.hash).await?;
983 let merge_base = nearest_ancestor_in(&source_git, &history, &side.shared).await?;
984 let mut divergence = Divergence {
985 head: head.hash.clone(),
986 base: base.hash.clone(),
987 merge_base: merge_base.clone(),
988 behind,
989 ..Divergence::default()
990 };
991 let merge_base_tree = match &merge_base {
992 Some(hash) => target_history
993 .iter()
994 .find(|commit| commit.hash == *hash)
995 .map(|commit| commit.tree_hash.clone()),
996 None => None,
997 };
998 let Some(merge_base_tree) = merge_base_tree else {
999 // No common history to compare from: say nothing is known.
1000 divergence.truncated = true;
1001 return Ok(Some(divergence));
1002 };
1003 let (ours, truncated_ours) =
1004 diff::changed_paths(&source_git, Some(&merge_base_tree), &head.tree_hash).await?;
1005 divergence.ours = ours;
1006 divergence.truncated = truncated_ours;
1007 if behind {
1008 let now = now_ms();
1009 let key = (target.id.clone(), merge_base_tree.clone(), base.tree_hash.clone());
1010 let (theirs, truncated_theirs) = match THEIRS.with(|memo| memo.borrow().get(&key, now)) {
1011 Some(kept) => kept,
1012 None => {
1013 let found = diff::changed_paths(&target_git, Some(&merge_base_tree), &base.tree_hash).await?;
1014 THEIRS.with(|memo| memo.borrow_mut().put(key, found.clone(), now));
1015 found
1016 }
1017 };
1018 divergence.theirs = theirs;
1019 divergence.truncated |= truncated_theirs;
1020 }
1021 Ok(Some(divergence))
1022 }
1023
1024 /// A target branch's history from its head, worked out once per head
1025 /// for every pull request asking about it (coalesce.rs). The head is
1026 /// read under the refs version; the history by its hash, which the
1027 /// object cache keeps for good.
1028 async fn target_side<R: GitRepo>(&self, target: &Repo, git: &R, branch: &str) -> Result<Rc<coalesce::TargetSide>> {
1029 let now = now_ms();
1030 let key = refs_cache::usable(registry::refs_state(&target.id), now)
1031 .map(|version| (target.id.clone(), branch.to_owned(), version));
1032 if let Some(key) = &key
1033 && let Some(side) = TARGETS.with(|memo| memo.borrow().get(key, now))
1034 {
1035 return Ok(side);
1036 }
1037 let history = match git.log(branch, 1).await?.first() {
1038 Some(head) => git.log(&head.hash, MAX_ANCESTRY).await?,
1039 None => Vec::new(),
1040 };
1041 let side = coalesce::TargetSide::new(history);
1042 if let Some(key) = key {
1043 TARGETS.with(|memo| memo.borrow_mut().put(key, side.clone(), now));
1044 }
1045 Ok(side)
1046 }
1047
1048 async fn head(&self, a: HeadArgs) -> Result<Option<String>> {
1049 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
1050 return Ok(None);
1051 };
1052 let branch = if a.branch.is_empty() { &repo.default_branch } else { &a.branch };
1053 let git = self.read_git(&repo).await?;
1054 Ok(git
1055 .log(branch, 1)
1056 .await?
1057 .into_iter()
1058 .next()
1059 .map(|commit| commit.hash))
1060 }
1061
1062 async fn delete_branch(&self, a: DeleteBranchArgs) -> Result<Outcome<bool>> {
1063 if !a.branch.starts_with(G1T_BRANCH_PREFIX) {
1064 return Ok(Outcome::fail(
1065 FailureCode::Forbidden,
1066 "Only branches g1t made for itself can be deleted this way.",
1067 ));
1068 }
1069 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
1070 return Ok(not_found());
1071 };
1072 let repo = match self.unpaused(repo).await? {
1073 Ok(repo) => repo,
1074 Err((code, message)) => return Ok(Outcome::fail(code, message)),
1075 };
1076 self.live(&repo).await?;
1077 let git = self.store.open(&store_key(&repo)).await?;
1078 let Some(old) = git
1079 .branches()
1080 .await?
1081 .into_iter()
1082 .find(|branch| branch.name == a.branch)
1083 .map(|branch| branch.hash)
1084 else {
1085 return Ok(Outcome::Ok(false));
1086 };
1087 let access = git.access(Scope::Write).await?;
1088 let deleted = land::delete_ref(&access, &a.branch, &old).await?;
1089 self.refs_moved(&repo.id).await;
1090 if let Err(reason) = deleted {
1091 return Ok(Outcome::fail(
1092 FailureCode::Conflict,
1093 format!("{} could not be deleted: {reason}", a.branch),
1094 ));
1095 }
1096 Ok(Outcome::Ok(true))
1097 }
1098
1099 async fn fork_for_pull(&self, a: ForkArgs) -> Result<Outcome<Repo>> {
1100 let viewer = Some(a.actor.clone());
1101 let Some(source) = self
1102 .registry
1103 .by_id(&a.source_id)
1104 .await?
1105 .filter(|repo| can_read(repo, &viewer))
1106 else {
1107 return Ok(not_found());
1108 };
1109 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1110 return Ok(Outcome::fail(code, message));
1111 }
1112 // Its working copy is made in its namespace: not while it moves.
1113 let source = match self.unpaused(source).await? {
1114 Ok(source) => source,
1115 Err((code, message)) => return Ok(Outcome::fail(code, message)),
1116 };
1117 let now = now_ms();
1118 let fork = Repo {
1119 id: new_id("rep", now),
1120 namespace: PULLS_NAMESPACE.to_owned(),
1121 name: a.pull_id.clone(),
1122 description: None,
1123 // A fork is exactly as visible as the repo it came from.
1124 is_private: source.is_private,
1125 owner_id: a.actor.id.clone(),
1126 default_branch: source.default_branch.clone(),
1127 fork_of: Some(source.id.clone()),
1128 protected: false,
1129 created_at: rfc3339(now),
1130 topics: Vec::new(),
1131 website: None,
1132 archived_at: None,
1133 };
1134 // Artifacts forks within a namespace: the copy goes where its
1135 // repository is.
1136 let (namespace, _) = store::locate(&store_key(&source));
1137 self.registry
1138 .claim_store_key(&fork, Some(&namespace), &self.store.default_namespace())
1139 .await?;
1140 self.store
1141 .open(&store_key(&source))
1142 .await?
1143 .fork(&store_key(&fork))
1144 .await?;
1145 self.registry.insert(&fork).await?;
1146 self.publish(NewEvent {
1147 kind: "repo.forked",
1148 source: SOURCE,
1149 repo_id: Some(source.id.clone()),
1150 actor: Some(a.actor.id),
1151 data: RepoForked {
1152 repo_id: fork.id.clone(),
1153 source_repo_id: source.id,
1154 pull_id: a.pull_id,
1155 },
1156 })
1157 .await?;
1158 Ok(Outcome::Ok(fork))
1159 }
1160
1161 async fn git_access(&self, a: GitAccessArgs) -> Result<Outcome<GitAccess>> {
1162 let found = self.registry.by_path(&a.path).await?;
1163 Ok(match self.authorize_git(&a.path, &a.viewer, a.service, found).await? {
1164 Outcome::Ok(repo) => {
1165 self.live(&repo).await?;
1166 let write = a.service == GitService::ReceivePack;
1167 if write {
1168 // A push with this credential would not pass through
1169 // here, so nothing that lists the refs is kept until it
1170 // has expired (see refs_cache.rs).
1171 let until = now_ms() + store::CREDENTIAL_LIFE_MS + 60_000;
1172 if let Err(error) = self.registry.refs_open(&repo.id, until).await {
1173 // Before the column exists nothing is kept anyway.
1174 if registry::refs_state(&repo.id).is_some() {
1175 return Err(error);
1176 }
1177 }
1178 }
1179 let scope = if write { Scope::Write } else { Scope::Read };
1180 Outcome::Ok(self.store.handout(&store_key(&repo), scope).await?)
1181 }
1182 Outcome::Fail(failure) => Outcome::Fail(failure),
1183 })
1184 }
1185
1186 /// The repository at `path` (`found`, as just read), if the viewer may
1187 /// use `service` on it: fetch from it, or push to it. A push to a path
1188 /// with nothing there makes the repository, in a workspace the pusher
1189 /// belongs to.
1190 async fn authorize_git(
1191 &self,
1192 path: &RepoPath,
1193 viewer: &Viewer,
1194 service: GitService,
1195 found: Option<Repo>,
1196 ) -> Result<Outcome<Repo>> {
1197 let mut a = GitAccessArgs {
1198 path: path.clone(),
1199 viewer: viewer.clone(),
1200 service,
1201 };
1202 let write = a.service == GitService::ReceivePack;
1203 // An access token: pushing needs code:write, reading a private
1204 // repository code:read. A public repository reads as it would for
1205 // anyone. Which repositories a token reaches is its owner's, checked
1206 // below as for anyone.
1207 if let Some(access) = a.viewer.as_ref().and_then(|user| user.token.as_deref()).cloned() {
1208 let public = found.as_ref().is_some_and(|repo| !repo.is_private);
1209 let decision = g1t_contracts::scopes::decide_git(&access, write, public);
1210 if !decision.allowed {
1211 return Ok(Outcome::fail(
1212 FailureCode::Forbidden,
1213 format!("{}\n", decision.reason.unwrap_or_default()),
1214 ));
1215 }
1216 if !write && !access.allows(g1t_contracts::scopes::Scope::CodeRead) {
1217 a.viewer = None;
1218 }
1219 }
1220
1221 // Anonymous callers are asked to authenticate whether or not the repo
1222 // exists, so private repos cannot be told apart from missing ones.
1223 let denied = || match &a.viewer {
1224 Some(_) => not_found(),
1225 None => Outcome::fail(FailureCode::Unauthenticated, "Authentication required."),
1226 };
1227 // An agent's token works through the API only: its sandbox has its
1228 // own way to push, to its own pull request.
1229 if a.viewer.as_ref().is_some_and(|user| user.kind == PrincipalKind::Agent) {
1230 return Ok(Outcome::fail(
1231 FailureCode::Forbidden,
1232 "A g1t agent's token cannot be used with git.",
1233 ));
1234 }
1235 if let (true, Some(user)) = (write, &a.viewer)
1236 && !user.verified
1237 {
1238 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1239 }
1240 let repo = match found {
1241 Some(repo) => {
1242 let allowed = if write {
1243 can_write(&repo, &a.viewer)
1244 } else {
1245 self.may_read(&repo, &a.viewer).await?
1246 };
1247 if !allowed {
1248 return Ok(denied());
1249 }
1250 // An archived repository, or a pull request's copy of one,
1251 // is read-only.
1252 if write {
1253 let archived = match &repo.fork_of {
1254 Some(source) => self.registry.by_id(source).await?,
1255 None => Some(repo.clone()),
1256 };
1257 match archived {
1258 Some(source) => {
1259 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1260 return Ok(Outcome::fail(code, format!("{message}\n")));
1261 }
1262 }
1263 // The repository it was copied from is deleted.
1264 None => return Ok(denied()),
1265 }
1266 }
1267 repo
1268 }
1269 None => {
1270 // Push to create, in a workspace the pusher belongs to.
1271 let owner = a
1272 .viewer
1273 .as_ref()
1274 .filter(|user| write && user.is_member(&a.path.namespace.to_lowercase()));
1275 let Some(owner) = owner else {
1276 return Ok(denied());
1277 };
1278 let created = self.create(push_to_create(owner, &a.path)).await?;
1279 match created {
1280 Outcome::Ok(repo) => repo,
1281 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
1282 }
1283 }
1284 };
1285 // A push, or a credential to push with, waits while the repository
1286 // moves between namespaces (moves.rs), and goes to where it is now.
1287 if write {
1288 return Ok(match self.unpaused(repo).await? {
1289 Ok(repo) => Outcome::Ok(repo),
1290 Err((code, message)) => Outcome::fail(code, format!("{message}\n")),
1291 });
1292 }
1293 Ok(Outcome::Ok(repo))
1294 }
1295
1296 async fn land(&self, a: LandArgs) -> Result<Outcome<Landed>> {
1297 let actor: Viewer = Some(a.actor.clone());
1298 let Some(source) = self.registry.by_id(&a.source_id).await? else {
1299 return Ok(not_found());
1300 };
1301 // A fork lands on the repository it came from; a branch on its own.
1302 let target = match &source.fork_of {
1303 Some(id) => self.registry.by_id(id).await?,
1304 None => Some(source.clone()),
1305 };
1306 let Some(target) = target.filter(|repo| can_read(repo, &actor)) else {
1307 return Ok(not_found());
1308 };
1309 if !registry::can(&target, &actor, Capability::Merge) {
1310 return Ok(Outcome::fail(
1311 FailureCode::Forbidden,
1312 access::needs(Capability::Merge, &format!("{}/{}", target.namespace, target.name)),
1313 ));
1314 }
1315 if !a.actor.verified {
1316 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1317 }
1318 if let Some((code, message)) = lifecycle::archived_refusal(&target) {
1319 return Ok(Outcome::fail(code, message));
1320 }
1321 // Moving between namespaces: wait for it (moves.rs). Both are read
1322 // again once it is done, for their new keys.
1323 let (source, target) = match (self.unpaused(source).await?, self.unpaused(target).await?) {
1324 (Ok(source), Ok(target)) => (source, target),
1325 (Err((code, message)), _) | (_, Err((code, message))) => return Ok(Outcome::fail(code, message)),
1326 };
1327
1328 let branch = &a.target_branch.clone().unwrap_or_else(|| target.default_branch.clone());
1329 let from_fork = source.id != target.id;
1330 let source_branch = match a.branch {
1331 Some(name) if !from_fork && name == *branch => {
1332 return Ok(Outcome::fail(
1333 FailureCode::Invalid,
1334 format!("{branch} cannot be merged into itself."),
1335 ));
1336 }
1337 Some(name) => name,
1338 None if from_fork => branch.clone(),
1339 None => {
1340 return Ok(Outcome::fail(
1341 FailureCode::Invalid,
1342 "Say which branch to merge.",
1343 ));
1344 }
1345 };
1346
1347 self.live(&source).await?;
1348 let source_git = self.store.open(&store_key(&source)).await?;
1349 let target_git = self.store.open(&store_key(&target)).await?;
1350 let history = source_git.log(&source_branch, MAX_ANCESTRY).await?;
1351 let Some(new) = history.first().map(|commit| commit.hash.clone()) else {
1352 return Ok(Outcome::fail(
1353 FailureCode::Conflict,
1354 "This pull request has no commits to merge.",
1355 ));
1356 };
1357 let old = target_git
1358 .log(branch, 1)
1359 .await?
1360 .into_iter()
1361 .next()
1362 .map(|commit| commit.hash);
1363
1364 if old.as_deref() == Some(new.as_str()) {
1365 return Ok(Outcome::Ok(Landed {
1366 commit: new,
1367 previous: None,
1368 }));
1369 }
1370 // Moving the branch to a commit that does not descend from its
1371 // current head would discard whatever landed in between.
1372 if let Some(old) = &old
1373 && !descends_from(&source_git, &history, old).await?
1374 {
1375 let remedy = if from_fork {
1376 format!("Pull {branch} into the pull request's fork, push, and merge again.")
1377 } else {
1378 format!("Merge {branch} into {source_branch}, push, and merge again.")
1379 };
1380 return Ok(Outcome::fail(
1381 FailureCode::Conflict,
1382 format!("{branch} has moved since this pull request was opened. {remedy}"),
1383 ));
1384 }
1385
1386 // For a branch the objects are already in the target; sending them
1387 // again is harmless and keeps one way of moving a ref.
1388 let source_access = source_git.access(Scope::Read).await?;
1389 let target_access = target_git.access(Scope::Write).await?;
1390 let pushed =
1391 land::fast_forward(&source_access, &target_access, branch, old.as_deref(), &new)
1392 .await?;
1393 self.refs_moved(&target.id).await;
1394 if let Err(reason) = pushed {
1395 // Most often another pull request landed between the check and the push.
1396 return Ok(Outcome::fail(
1397 FailureCode::Conflict,
1398 format!("{branch} could not be updated: {reason}"),
1399 ));
1400 }
1401 self.publish_push(
1402 &target,
1403 &format!("refs/heads/{branch}"),
1404 old.as_deref(),
1405 &new,
1406 Some(a.actor.id),
1407 )
1408 .await?;
1409 Ok(Outcome::Ok(Landed {
1410 commit: new,
1411 previous: old,
1412 }))
1413 }
1414
1415 async fn compare(&self, a: CompareArgs) -> Result<Outcome<Comparison>> {
1416 let Some(repo) = self
1417 .visible(self.registry.by_id(&a.repo_id).await?, &a.viewer)
1418 .await?
1419 else {
1420 return Ok(not_found());
1421 };
1422 let git = self.read_git(&repo).await?;
1423 let head_ref = a.head.as_deref().unwrap_or(&repo.default_branch);
1424 // A pull request into another branch is compared from where it
1425 // left that branch.
1426 let base_branch = a.base_branch.clone();
1427 // The head's history is only searched when the base is worked out
1428 // from another branch.
1429 let depth = if a.base.is_some() || is_commit_hash(head_ref) { 1 } else { MAX_ANCESTRY };
1430 let history = git.log(head_ref, depth).await?;
1431 let Some(head) = history.first() else {
1432 return Ok(Outcome::fail(
1433 FailureCode::Conflict,
1434 "There are no commits to compare.",
1435 ));
1436 };
1437
1438 // Where the head's history meets the default branch of `against`,
1439 // or the branch asked for.
1440 let shared_with = async |against: &Repo| -> Result<Option<String>> {
1441 let against_git = self.read_git(against).await?;
1442 let branch = base_branch.as_deref().unwrap_or(&against.default_branch);
1443 let shared: HashSet<String> = against_git
1444 .log(branch, MAX_ANCESTRY)
1445 .await?
1446 .into_iter()
1447 .map(|commit| commit.hash)
1448 .collect();
1449 nearest_ancestor_in(&git, &history, &shared).await
1450 };
1451 let base = match (a.base, &repo.fork_of) {
1452 (Some(base), _) => Some(base),
1453 // A fork is compared with the last commit it shares with the
1454 // repository it came from.
1455 (None, Some(target_id)) => match self.registry.by_id(target_id).await? {
1456 Some(target) => shared_with(&target).await?,
1457 None => None,
1458 },
1459 // A branch, with the point where it left the default branch.
1460 // A single commit, with its first parent.
1461 (None, None) if is_commit_hash(head_ref) => head.parents.first().cloned(),
1462 (None, None) if head_ref != base_branch.as_deref().unwrap_or(&repo.default_branch) => {
1463 shared_with(&repo).await?
1464 }
1465 (None, None) => head.parents.first().cloned(),
1466 };
1467 let base_tree = match &base {
1468 Some(base) => git
1469 .log(base, 1)
1470 .await?
1471 .into_iter()
1472 .next()
1473 .map(|commit| commit.tree_hash),
1474 None => None,
1475 };
1476 let (files, truncated) =
1477 diff::compare_trees(&git, base_tree.as_deref(), &head.tree_hash).await?;
1478 Ok(Outcome::Ok(Comparison {
1479 base,
1480 head: head.hash.clone(),
1481 files,
1482 truncated,
1483 }))
1484 }
1485
1486 /// Reports that `git_ref` of `repo` (a full ref) now points to `after`.
1487 async fn publish_push(
1488 &self,
1489 repo: &Repo,
1490 git_ref: &str,
1491 before: Option<&str>,
1492 after: &str,
1493 actor: Option<String>,
1494 ) -> Result<()> {
1495 self.publish_git_push(repo, git_ref, before, after, actor, false).await
1496 }
1497
1498 /// `publish_push`, saying whether the push reached the store without
1499 /// being scanned for secrets first.
1500 async fn publish_git_push(
1501 &self,
1502 repo: &Repo,
1503 git_ref: &str,
1504 before: Option<&str>,
1505 after: &str,
1506 actor: Option<String>,
1507 unscanned: bool,
1508 ) -> Result<()> {
1509 self.publish(NewEvent {
1510 kind: "git.push",
1511 source: SOURCE,
1512 repo_id: Some(repo.id.clone()),
1513 actor,
1514 data: GitPush {
1515 repo_id: repo.id.clone(),
1516 git_ref: git_ref.to_owned(),
1517 before: before.map(str::to_owned),
1518 after: after.to_owned(),
1519 default_branch: git_ref.strip_prefix("refs/heads/")
1520 == Some(repo.default_branch.as_str()),
1521 unscanned,
1522 },
1523 })
1524 .await
1525 }
1526
1527 /// Git over HTTPS. Only what decides the answer happens before it:
1528 /// the repository, who is asking and whether they may, the free
1529 /// workspace limits, push protection, and the store's own answer. The
1530 /// audit entry and what a push changed are recorded once git has its
1531 /// answer. Each answer says how long its steps took (`Server-Timing`).
1532 async fn git_http(&self, request: Request, env: &Env, ctx: &Context) -> Result<Response> {
1533 let mut timing = git_http::Timing::start();
1534 let Some(git) = git_http::parse(&request.url()?) else {
1535 return Response::error("Not found", 404);
1536 };
1537 let response = match self.answer_git(request, &git, env, ctx, &mut timing).await {
1538 Ok(response) => response,
1539 // The git store is busy: git hears when to try again.
1540 Err(error) => match resilience::busy(&error.to_string()) {
1541 Some(busy) => git_http::busy_response(busy)?,
1542 None => return Err(error),
1543 },
1544 };
1545 timing.apply(response)
1546 }
1547
1548 async fn answer_git(
1549 &self,
1550 request: Request,
1551 git: &git_http::GitRequest,
1552 env: &Env,
1553 ctx: &Context,
1554 timing: &mut git_http::Timing,
1555 ) -> Result<Response> {
1556 let write = git.service == GitService::ReceivePack;
1557 let get = request.method() == Method::Get;
1558 let identity = env.service("IDENTITY")?;
1559 // The repository and the caller's credentials, at once. A fetch may
1560 // go by the row as read a moment ago, for the same clone's next
1561 // request; a push always reads it. Anonymous callers cost nothing.
1562 let lookup = async {
1563 if write {
1564 self.registry.by_path(&git.path).await
1565 } else {
1566 self.registry.by_path_recent(&git.path).await
1567 }
1568 };
1569 let (found, viewer) =
1570 futures_util::future::join(lookup, git_http::viewer(&request, &identity)).await;
1571 let found = found?;
1572 timing.mark("repo");
1573 if found.is_none() {
1574 // A workspace that was renamed: git follows a redirect when it
1575 // first asks for refs, and uses the new address from then on.
1576 // A repository transferred to another workspace: the same, to
1577 // its new path. Fetches and pushes both follow either.
1578 let url = request.url()?;
1579 let (renamed, moved) = futures_util::future::join(
1580 git_http::renamed(&url, &identity),
1581 self.registry.resolve_moved(&git.path),
1582 )
1583 .await;
1584 timing.mark("moved");
1585 if let Some(location) = renamed? {
1586 return git_http::moved(&location, get);
1587 }
1588 if let Some(now) = moved?
1589 && let Some(location) = git_http::transferred(&url, &now)
1590 {
1591 return git_http::moved(&location, get);
1592 }
1593 }
1594 let viewer = viewer?;
1595 // A run credential is checked against its grants, then acts as the
1596 // person it works for. See run_access.rs.
1597 let (request, viewer, audit) = match self.admit_git(request, git, viewer, found.as_ref()).await? {
1598 run_access::Admitted::Go { request, viewer, entry } => (request, viewer, entry),
1599 run_access::Admitted::Refused(response) => return Ok(response),
1600 };
1601 let mut after = AfterGit {
1602 audit,
1603 status: 0,
1604 message: None,
1605 push: None,
1606 };
1607 let repo = match self.authorize_git(&git.path, &viewer, git.service, found).await? {
1608 Outcome::Ok(repo) => repo,
1609 refused => {
1610 let response = git_http::refuse(refused)?;
1611 after.ended(response.status_code(), None);
1612 after.spawn(env, ctx);
1613 return Ok(response);
1614 }
1615 };
1616 // A pull request's working copy removed after it closed is made
1617 // again before git uses it (forks.rs).
1618 self.live(&repo).await?;
1619 timing.mark("access");
1620 // A protected default branch takes changes only from a merged pull
1621 // request, which lands without going through here.
1622 let protected = (repo.protected && repo.fork_of.is_none()).then(|| repo.default_branch.clone());
1623 // Clones check out the default branch g1t keeps, which can have
1624 // changed since the store made the repository.
1625 let default_branch = repo.fork_of.is_none().then(|| repo.default_branch.clone());
1626 let key = store_key(&repo);
1627 let scope = if write { Scope::Write } else { Scope::Read };
1628 let mut request = request;
1629 let protocol = refs_cache::protocol(request.headers().get("git-protocol")?.as_deref());
1630 // A fetch's POST is read here, to tell an `ls-refs` from a fetch of
1631 // objects; the store would have it read in full anyway.
1632 let body = if !write && !get { Some(request.bytes().await?) } else { None };
1633 // What it asks the store, for the meters (meters.rs).
1634 let call = git_ops::classify(git.service, git.endpoint, get, body.as_deref());
1635 // Answers kept from the usual store may name refs the fallback
1636 // store does not have (fallback.rs): none are used, or kept.
1637 let fallback = self.store.on_fallback(&key);
1638 // An answer that lists refs may have been kept: see refs_cache.rs.
1639 let kept_key = refs_cache::kind(git, get, protocol, body.as_deref())
1640 .filter(|_| !fallback)
1641 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1642 .map(|(kind, version)| {
1643 refs_cache::Key::new(&repo.id, version, default_branch.as_deref(), protocol, &kind)
1644 });
1645 // A fresh clone's pack may have been kept too: see pack_cache.rs.
1646 // Under the same refs version, so never across a change to them.
1647 let pack_key = self
1648 .packs
1649 .as_ref()
1650 .filter(|_| !fallback)
1651 .and_then(|_| {
1652 let encoding = request.headers().get("content-encoding").ok().flatten();
1653 pack_cache::cacheable(git, get, protocol, encoding.as_deref(), body.as_deref())
1654 })
1655 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1656 .map(|(normalized, version)| pack_cache::Key::new(&repo.id, version, &normalized));
1657 // A kept answer and the free workspace limits, with a kept
1658 // credential looked up alongside. A kept answer goes back without
1659 // waiting for the credential, which it does not need.
1660 let ((answer, pack, limited), kept_access) = {
1661 let shared = self.shared.as_deref();
1662 let answer_and_limits = std::pin::pin!(futures_util::future::join3(
1663 async {
1664 match &kept_key {
1665 Some(kept_key) => refs_cache::get(shared, kept_key).await,
1666 None => None,
1667 }
1668 },
1669 async {
1670 match (&pack_key, self.packs.as_deref()) {
1671 (Some(pack_key), Some(packs)) => pack_cache::get(packs, pack_key).await,
1672 _ => None,
1673 }
1674 },
1675 self.git_limits(call, git, &repo, env),
1676 ));
1677 let kept_access = std::pin::pin!(self.store.kept_access(&key, scope));
1678 match futures_util::future::select(answer_and_limits, kept_access).await {
1679 futures_util::future::Either::Left((first, kept_access)) => {
1680 let answered = first.0.is_some() || first.1.is_some() || matches!(first.2, Ok(Some(_)) | Err(_));
1681 (first, if answered { None } else { kept_access.await })
1682 }
1683 futures_util::future::Either::Right((kept_access, first)) => (first.await, kept_access),
1684 }
1685 };
1686 timing.mark("kept");
1687 // A kept pack first: it never reaches the store, so it is never an
1688 // operation, and a free workspace past its operation cap still gets
1689 // it. (The other limits are a push's, and a pack is only a fetch.)
1690 if let Some(kept) = pack {
1691 timing.note("pack", "hit");
1692 let sent = body.as_ref().map_or(0, |body| body.len() as u64);
1693 meters::record(pack_cache::HIT, &key, sent, kept.size);
1694 after.ended(200, None);
1695 after.spawn(env, ctx);
1696 return kept.response();
1697 }
1698 if let Some((response, status, message)) = limited? {
1699 after.ended(status, Some(message.to_owned()));
1700 after.spawn(env, ctx);
1701 return Ok(response);
1702 }
1703 if pack_key.is_some() {
1704 timing.note("pack", "miss");
1705 }
1706 if let (Some((entry, found)), Some(kept_key)) = (answer, &kept_key) {
1707 timing.note("refs", found.as_str());
1708 if found == refs_cache::Found::Shared {
1709 let (kept_key, entry) = (kept_key.clone(), entry.clone());
1710 ctx.wait_until(async move { refs_cache::keep_in_colo(&kept_key, &entry).await });
1711 }
1712 // Never reached the store: never an operation.
1713 meters::record(call.cached_meter(), &key, 0, entry.body.len() as u64);
1714 after.ended(200, None);
1715 after.spawn(env, ctx);
1716 return entry.response();
1717 }
1718 if kept_key.is_some() {
1719 timing.note("refs", "miss");
1720 }
1721 // The store's credential: one made a moment ago, here or in another
1722 // isolate (see store.rs), or a new one.
1723 let access = match kept_access {
1724 Some((access, from)) => {
1725 timing.note("cred", from.as_str());
1726 access
1727 }
1728 None => {
1729 let access = self.store.mint_access(&key, scope).await?;
1730 timing.mark("mint");
1731 timing.note("cred", "mint");
1732 access
1733 }
1734 };
1735 // Should the store turn a kept credential down, a fetch's first
1736 // request is tried again with a new one; the requests after it then
1737 // have that one too.
1738 let again = if get { Some(request.clone()?) } else { None };
1739 // Push protection: a push that adds a secret is refused. See secret_scan.rs.
1740 let scan = async |body: &[u8]| self.protect(&repo, viewer.as_ref(), body).await;
1741 // What a push may bring (pack_limits.rs): the repository's size is
1742 // its own and its pull requests' working copies'.
1743 let limits = if write && !get {
1744 git_http::PushLimits {
1745 held: self.held(&repo).await,
1746 repo_limit: self.repo_limit,
1747 large: self.large_pushes,
1748 ..git_http::PushLimits::default()
1749 }
1750 } else {
1751 git_http::PushLimits::default()
1752 };
1753 let mut outcome = git_http::forward(
1754 request,
1755 body,
1756 git,
1757 &access,
1758 protected.as_deref(),
1759 default_branch.as_deref(),
1760 limits,
1761 scan,
1762 )
1763 .await?;
1764 let turned_down = matches!(
1765 &outcome,
1766 git_http::Push::Forwarded(forwarded) if matches!(forwarded.response.status_code(), 401 | 403)
1767 );
1768 if turned_down {
1769 self.store.forget_access(&key).await;
1770 if let Some(again) = again {
1771 let access = self.store.mint_access(&key, scope).await?;
1772 let nothing = async |_: &[u8]| Ok(None);
1773 outcome = git_http::forward(
1774 again,
1775 None,
1776 git,
1777 &access,
1778 protected.as_deref(),
1779 default_branch.as_deref(),
1780 git_http::PushLimits::default(),
1781 nothing,
1782 )
1783 .await?;
1784 }
1785 }
1786 let forwarded =
1787 match outcome {
1788 git_http::Push::Forwarded(forwarded) => forwarded,
1789 git_http::Push::Refused(response) => {
1790 after.ended(403, Some("The push would change a protected branch.".to_owned()));
1791 after.spawn(env, ctx);
1792 return Ok(response);
1793 }
1794 git_http::Push::Blocked(response) => {
1795 after.ended(403, Some("The push adds a secret.".to_owned()));
1796 after.spawn(env, ctx);
1797 return Ok(response);
1798 }
1799 git_http::Push::Declined(response, reason) => {
1800 after.ended(403, Some(format!("The push was declined: {reason}.")));
1801 after.spawn(env, ctx);
1802 return Ok(response);
1803 }
1804 };
1805 if forwarded.from_store {
1806 let received = forwarded
1807 .response
1808 .headers()
1809 .get("content-length")?
1810 .and_then(|length| length.parse().ok())
1811 .unwrap_or(0);
1812 meters::record(call.meter(), &key, forwarded.sent, received);
1813 }
1814 timing.mark("store");
1815 let mut response = forwarded.response;
1816 let status = response.status_code();
1817 if write && !get {
1818 // A push: the store has moved its refs once it has answered in
1819 // full, so the answer is read before the change is recorded, and
1820 // only then goes back. Whoever fetches after it sees the push.
1821 let headers = response.headers().clone();
1822 headers.delete("content-length")?;
1823 let report = response.bytes().await?;
1824 self.refs_moved(&repo.id).await;
1825 timing.mark("refs");
1826 response = Response::from_bytes(report)?.with_headers(headers).with_status(status);
1827 } else if let (Some(kept_key), 200) = (&kept_key, status) {
1828 // A miss: this answer is kept for the next to ask.
1829 let headers = response.headers().clone();
1830 headers.delete("content-length")?;
1831 let body = response.bytes().await?;
1832 if let Some(content_type) = headers.get("content-type")? {
1833 let entry = refs_cache::Entry { content_type, body: body.clone() };
1834 if entry.keepable() {
1835 let shared = self.shared.clone();
1836 let kept_key = kept_key.clone();
1837 ctx.wait_until(async move { refs_cache::keep(shared.as_deref(), &kept_key, &entry).await });
1838 }
1839 }
1840 response = Response::from_bytes(body)?.with_headers(headers).with_status(status);
1841 } else if let (Some(pack_key), Some(packs), true) = (&pack_key, &self.packs, forwarded.from_store) {
1842 // A fresh clone the bucket did not have: counted, and its pack
1843 // kept as it streams to git, when it is a whole one.
1844 meters::record(pack_cache::MISS, &key, forwarded.sent, 0);
1845 if status == 200 {
1846 let store_key = key.clone();
1847 let measured = Box::new(move |bytes: u64| meters::record_bytes(pack_cache::MISS, &store_key, 0, bytes));
1848 let (teed, filling) = pack_cache::tee(response, packs.clone(), pack_key, measured)?;
1849 response = teed;
1850 if let Some(filling) = filling {
1851 let pack_key = pack_key.clone();
1852 ctx.wait_until(async move {
1853 let filled = filling.await;
1854 if !matches!(filled, pack_cache::Filled::Kept { .. } | pack_cache::Filled::Abandoned) {
1855 worker::console_warn!("pack {} not kept: {filled:?}", pack_key.as_str());
1856 }
1857 });
1858 }
1859 }
1860 }
1861 after.ended(status, None);
1862 if status == 200 && (forwarded.pack_bytes > 0 || !forwarded.pushed.is_empty()) {
1863 after.push = Some(PushDone {
1864 repo,
1865 pushed: forwarded.pushed,
1866 pack_bytes: forwarded.pack_bytes,
1867 actor: viewer.map(|user: User| user.id),
1868 unscanned: forwarded.unscanned,
1869 });
1870 }
1871 after.spawn(env, ctx);
1872 Ok(response)
1873 }
1874
1875 /// The answer for a request a free workspace's limits stop, or a push
1876 /// to a full repository, with its status and reason for the audit log;
1877 /// `None` to go on.
1878 ///
1879 /// A clone, fetch or push is a git operation, which the git store
1880 /// charges g1t for: counted for billing once the answer has gone back
1881 /// (meters.rs), and a free workspace far past its share is slowed down
1882 /// rather than charged (see git_ops.rs). Whether it is past it is
1883 /// decided from counts this isolate already holds: the database is not
1884 /// asked on the way. A free workspace is never charged for private
1885 /// storage: once its private repositories hold the free amount, pushes
1886 /// to them stop, checked when a push begins so that git shows the
1887 /// reason. So do pushes to a repository at the store's size limit.
1888 async fn git_limits(
1889 &self,
1890 call: git_ops::GitCall,
1891 git: &git_http::GitRequest,
1892 repo: &Repo,
1893 env: &Env,
1894 ) -> Result<Option<(Response, u16, &'static str)>> {
1895 let namespace = git.path.namespace.to_lowercase();
1896 if meters::mapping_now().billable(call.meter()) > 0.0 {
1897 let now = now_ms();
1898 let hour = git_ops::hour_key(&rfc3339(now));
1899 let limits = git_ops::Limits::from_env(env);
1900 if let Some((month, hour_ops)) = git_ops::standing(&namespace, &hour, now)
1901 && git_ops::slow_down(month + 1, hour_ops + 1, limits.free_cap, limits.hourly)
1902 && git_ops::is_free_kept(env.service("BILLING").ok().as_ref(), &namespace).await
1903 {
1904 return Ok(Some((
1905 git_ops::too_many(&namespace, limits.free_cap, limits.hourly)?,
1906 429,
1907 "Too many git operations this hour.",
1908 )));
1909 }
1910 }
1911 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" {
1912 let held = self.held(repo).await;
1913 if held >= self.repo_limit {
1914 let message = format!(
1915 "{}/{} holds about {}, the most a repository may hold on g1t, so it takes no more pushes. Delete what you no longer need, or split it: https://docs.g1t.sh/guides/git/#size-limits\n",
1916 repo.namespace,
1917 repo.name,
1918 pack_limits::megabytes(held)
1919 );
1920 return Ok(Some((Response::error(message, 403)?, 403, "The repository is full.")));
1921 }
1922 }
1923 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" && repo.is_private {
1924 let free = git_ops::free_private_bytes(env);
1925 let held = self.registry.private_bytes(&namespace).await.unwrap_or(0);
1926 if git_ops::storage_full(held, free)
1927 && git_ops::is_free(env.service("BILLING").ok().as_ref(), &namespace).await
1928 {
1929 return Ok(Some((
1930 git_ops::storage_full_response(&namespace, held, free)?,
1931 403,
1932 "Free private storage is full.",
1933 )));
1934 }
1935 }
1936 Ok(None)
1937 }
1938
1939 /// What a repository and its pull requests' working copies hold, as
1940 /// g1t counts it: read for a push's first request, kept a minute for
1941 /// the rest of it.
1942 async fn held(&self, repo: &Repo) -> u64 {
1943 let root = repo.fork_of.clone().unwrap_or_else(|| repo.id.clone());
1944 let now = now_ms();
1945 if let Some(held) = HELD.with(|held| held.borrow().get(&root, now)) {
1946 return held;
1947 }
1948 let held = self.registry.stored_bytes(&root).await.unwrap_or(0).max(0) as u64;
1949 HELD.with(|kept| kept.borrow_mut().put(root, held, now));
1950 held
1951 }
1952
1953 /// What a push changed, recorded once git has its answer.
1954 async fn record_push(&self, push: PushDone) -> Result<()> {
1955 let PushDone {
1956 repo,
1957 pushed,
1958 pack_bytes,
1959 actor,
1960 unscanned,
1961 } = push;
1962 // What the push stored, for billing's storage meter. A failure only
1963 // leaves the count short.
1964 if pack_bytes > 0
1965 && let Err(error) = self.registry.add_stored_bytes(&repo, pack_bytes).await
1966 {
1967 worker::console_error!("stored bytes for {} not counted: {error}", repo.name);
1968 }
1969 if pushed.is_empty() {
1970 return Ok(());
1971 }
1972 // Artifacts' own push notifications are per repository, which does
1973 // not fit a repo per pull request, so the front end reports pushes
1974 // itself: one event for each branch that moved.
1975 let stored = self.store.open(&store_key(&repo)).await?;
1976 for pushed in &pushed {
1977 // The store can refuse one ref and accept another, so each
1978 // branch is checked against where it actually is. A tag the
1979 // store cannot read back is taken as pushed.
1980 let moved = match pushed.branch() {
1981 Some(branch) => stored
1982 .log(branch, 1)
1983 .await?
1984 .first()
1985 .is_some_and(|commit| commit.hash == pushed.after),
1986 None => stored.log(&pushed.git_ref, 1).await.map_or(true, |head| {
1987 head.first().is_none_or(|commit| commit.hash == pushed.after)
1988 }),
1989 };
1990 if moved {
1991 self.publish_git_push(
1992 &repo,
1993 &pushed.git_ref,
1994 pushed.before.as_deref(),
1995 &pushed.after,
1996 actor.clone(),
1997 unscanned,
1998 )
1999 .await?;
2000 }
2001 }
2002 Ok(())
2003 }
2004}
2005
2006/// A push the store accepted, to be recorded once git has its answer.
2007struct PushDone {
2008 repo: Repo,
2009 pushed: Vec<git_http::Pushed>,
2010 pack_bytes: u64,
2011 actor: Option<String>,
2012 /// Too large to scan for secrets before it was stored.
2013 unscanned: bool,
2014}
2015
2016/// What a git request leaves for after its answer: its audit entry, with
2017/// how the request ended, and what a push changed.
2018struct AfterGit {
2019 audit: Option<Box<g1t_contracts::audit::NewAuditEntry>>,
2020 status: u16,
2021 message: Option<String>,
2022 push: Option<PushDone>,
2023}
2024
2025impl AfterGit {
2026 fn ended(&mut self, status: u16, message: Option<String>) {
2027 self.status = status;
2028 self.message = message;
2029 }
2030
2031 /// Does the work once the response is on its way. A failure is logged:
2032 /// git has already been told how its request went.
2033 fn spawn(self, env: &Env, ctx: &Context) {
2034 if self.audit.is_none() && self.push.is_none() {
2035 return;
2036 }
2037 let env = env.clone();
2038 ctx.wait_until(async move {
2039 let repos = match service(&env) {
2040 Ok(repos) => repos,
2041 Err(error) => {
2042 worker::console_error!("git request not recorded: {error}");
2043 return;
2044 }
2045 };
2046 repos.finish_git(self.audit, self.status, self.message).await;
2047 if let Some(push) = self.push
2048 && let Err(error) = repos.record_push(push).await
2049 {
2050 worker::console_error!("push not recorded: {error}");
2051 }
2052 });
2053 }
2054}
2055
2056fn service(env: &Env) -> Result<Repos<ArtifactsStore>> {
2057 let shared = shared::Shared::from_env(env).map(Rc::new);
2058 Ok(Repos {
2059 registry: Registry { db: env.d1("DB")? },
2060 store: ArtifactsStore::new(env, shared.clone())?,
2061 shared,
2062 packs: pack_cache::Packs::from_env(env).map(Rc::new),
2063 events: env.service("EVENTS")?,
2064 security: env.service("SECURITY").ok(),
2065 billing: env.service("BILLING").ok(),
2066 identity: env.service("IDENTITY").ok(),
2067 free_private_bytes: git_ops::free_private_bytes(env),
2068 fork_days: forks::retention_days(env),
2069 repo_limit: env
2070 .var("REPO_STORAGE_LIMIT_BYTES")
2071 .ok()
2072 .and_then(|value| value.to_string().parse().ok())
2073 .unwrap_or(pack_limits::DEFAULT_REPO_LIMIT_BYTES),
2074 large_pushes: git_http::LargePushes::from_var(env.var("LARGE_PUSHES").ok().map(|value| value.to_string()).as_deref()),
2075 placement: shards::Placement::from_vars(
2076 env.var("ARTIFACTS_NEW_REPOS").ok().map(|value| value.to_string()).as_deref(),
2077 env.var("ARTIFACTS_EU_NAMESPACE").ok().map(|value| value.to_string()).as_deref(),
2078 ),
2079 limits: shards::limits(env.var("ARTIFACTS_NAMESPACE_LIMITS").ok().map(|value| value.to_string()).as_deref()),
2080 })
2081}
2082
2083/// Writes what this isolate metered once the answer has gone back, every
2084/// few seconds at most: now, or once it is due, waiting in this request's
2085/// `wait_until` so nothing counted is left for a request that may never
2086/// come (meters.rs).
2087fn flush_later(env: &Env, ctx: &Context) {
2088 let Some(wait) = meters::plan_flush() else {
2089 return;
2090 };
2091 if let Ok(db) = env.d1("DB") {
2092 ctx.wait_until(async move { meters::flush_after(&db, wait).await });
2093 }
2094}
2095
2096/// `/backups/<job id>/parts/<number>`: the job and the part's number.
2097fn backup_part_path(path: &str) -> Option<(String, u16)> {
2098 let rest = path.strip_prefix("/backups/")?;
2099 let (job, number) = rest.split_once("/parts/")?;
2100 let number = number.parse::<u16>().ok()?;
2101 (!job.is_empty() && !job.contains('/')).then(|| (job.to_owned(), number))
2102}
2103
2104fn backups_off<T>() -> Outcome<T> {
2105 Outcome::fail(FailureCode::Conflict, "Backups are off on this installation: it has no storage for them.")
2106}
2107
2108/// One part of a backup's bundle, with the job's token in its header.
2109async fn backup_part(request: &mut Request, env: &Env, repos: &Repos<ArtifactsStore>, job_id: String, number: u16) -> Result<Response> {
2110 let Some(blobs) = backups::storage(env) else {
2111 return reply(&backups_off::<()>());
2112 };
2113 let token = request.headers().get(g1t_contracts::backups::TOKEN_HEADER)?.unwrap_or_default();
2114 let bytes = request.bytes().await?;
2115 let job = g1t_contracts::backups::BackupJobArgs { job_id, token };
2116 reply(&backups::part(&repos.registry.db, &blobs, &job, number, bytes).await?)
2117}
2118
2119#[cfg(test)]
2120mod backup_path_tests {
2121 use super::backup_part_path;
2122
2123 #[test]
2124 fn a_part_is_named_by_its_job_and_number() {
2125 assert_eq!(backup_part_path("/backups/bkp_1/parts/3"), Some(("bkp_1".to_owned(), 3)));
2126 assert_eq!(backup_part_path("/backups/bkp_1/parts/x"), None);
2127 assert_eq!(backup_part_path("/backups//parts/1"), None);
2128 assert_eq!(backup_part_path("/acme/rocket.git/info/refs"), None);
2129 }
2130}
2131
2132/// Read methods whose answer is an `Outcome`: when the git store is busy,
2133/// the site is told so in words instead of failing the page.
2134const OUTCOME_READS: [&str; 6] = ["tree", "blob", "log", "branches", "blame", "compare"];
2135
2136#[event(fetch)]
2137async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
2138 let mut repos = service(&env)?;
2139 // A part of a backup's bundle, as the API passes it on from the
2140 // sandbox: bytes, not JSON (backups.rs).
2141 if request.method() == Method::Put
2142 && let Some((job_id, number)) = backup_part_path(&request.path())
2143 {
2144 let answered = backup_part(&mut request, &env, &repos, job_id, number).await;
2145 flush_later(&env, &ctx);
2146 return answered;
2147 }
2148 let Some(method) = rpc_method(&request) else {
2149 let answered = repos.git_http(request, &env, &ctx).await;
2150 flush_later(&env, &ctx);
2151 return answered;
2152 };
2153 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2154 // Git over HTTPS above always reads the primary.
2155 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2156 repos.registry.db = db;
2157 let body: serde_json::Value = request.json().await?;
2158
2159 let answered = async { match method.as_str() {
2160 "get" => reply(&repos.get(args(body)?).await?),
2161 "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
2162 "readable" => {
2163 let a: ReadableArgs = args(body)?;
2164 reply(&repos.registry.readable(&a.ids, &a.viewer).await?)
2165 }
2166 "public_namespaces" => {
2167 let a: PublicNamespacesArgs = args(body)?;
2168 reply(&repos.registry.public_namespaces(&a.owner_id).await?)
2169 }
2170 "path_by_id" => {
2171 let a: PathByIdArgs = args(body)?;
2172 reply(
2173 &repos
2174 .registry
2175 .by_id(&a.id)
2176 .await?
2177 .filter(|repo| repo.fork_of.is_none())
2178 .map(|repo| RepoPath {
2179 namespace: repo.namespace,
2180 name: repo.name,
2181 }),
2182 )
2183 }
2184 "list" => {
2185 let a: ListArgs = args(body)?;
2186 reply(
2187 &repos
2188 .registry
2189 .list(
2190 &a.viewer,
2191 a.query.as_deref(),
2192 a.namespace.as_deref(),
2193 a.member_only,
2194 )
2195 .await?,
2196 )
2197 }
2198 "create" => reply(&repos.create(args(body)?).await?),
2199 // Services only: a GitHub mirror catching up, or pushing out.
2200 "mirror" => reply(&repos.mirror(args(body)?).await?),
2201 "transfer" => reply(&repos.transfer(args(body)?).await?),
2202 // A repository's lifecycle: see lifecycle.rs.
2203 "delete" => reply(&repos.delete(args(body)?).await?),
2204 "deleted" => reply(&repos.deleted(args(body)?).await?),
2205 "restore" => reply(&repos.restore(args(body)?).await?),
2206 "purge" => reply(&repos.purge(args(body)?).await?),
2207 "purge_due" => reply(&repos.purge_due(args(body)?).await?),
2208 "rename" => reply(&repos.rename(args(body)?).await?),
2209 "archive" => reply(&repos.archive(args(body)?).await?),
2210 "set_visibility" => reply(&repos.set_visibility(args(body)?).await?),
2211 "set_default_branch" => reply(&repos.set_default_branch(args(body)?).await?),
2212 "rename_branch" => reply(&repos.rename_branch(args(body)?).await?),
2213 "resolve_branch" => reply(&repos.resolve_branch(args(body)?).await?),
2214 "status_by_id" => reply(&repos.status_by_id(args(body)?).await?),
2215 "resolve_path" => {
2216 let a: ResolvePathArgs = args(body)?;
2217 reply(&repos.registry.resolve_moved(&a.path).await?)
2218 }
2219 "namespace_count" => {
2220 let a: NamespaceCountArgs = args(body)?;
2221 reply(&repos.registry.count_in(&a.namespace).await?)
2222 }
2223 "update" => reply(&repos.update(args(body)?).await?),
2224 "tree" => reply(&repos.tree(args(body)?).await?),
2225 "blob" => reply(&repos.blob(args(body)?).await?),
2226 "log" => reply(&repos.log(args(body)?).await?),
2227 "blame" => reply(&repos.blame(args(body)?).await?),
2228 "fork_for_pull" => reply(&repos.fork_for_pull(args(body)?).await?),
2229 "git_access" => reply(&repos.git_access(args(body)?).await?),
2230 "branches" => reply(&repos.branches(args(body)?).await?),
2231 "last_commits" => reply(&repos.last_commits(args(body)?).await?),
2232 "tags" => reply(&repos.tags(args(body)?).await?),
2233 "head" => reply(&repos.head(args(body)?).await?),
2234 "behind" => reply(&repos.behind(args(body)?).await?),
2235 "divergence" => reply(&repos.divergence(args(body)?).await?),
2236 "land" => reply(&repos.land(args(body)?).await?),
2237 "update_pull_branch" => reply(&repos.update_pull_branch(args(body)?).await?),
2238 "delete_branch" => reply(&repos.delete_branch(args(body)?).await?),
2239 "commit_file" => reply(&repos.commit_file(args(body)?).await?),
2240 "compare" => reply(&repos.compare(args(body)?).await?),
2241 "scan_history" => reply(&repos.scan_history(args(body)?).await?),
2242 "find_lockfiles" => reply(&repos.find_lockfiles(args(body)?).await?),
2243 "match_pattern" => reply(&repos.match_pattern(args(body)?).await?),
2244 "check_secret" => reply(&repos.check_secret(args(body)?).await?),
2245 "list_files" => reply(&repos.list_files(args(body)?).await?),
2246 "changed_files" => reply(&repos.changed_files(args(body)?).await?),
2247 "read_blobs" => reply(&repos.read_blobs(args(body)?).await?),
2248 // Services only: what the Composer registry builds packages from.
2249 "refs" => reply(&repos.refs_of(args(body)?).await?),
2250 "raw_file" => reply(&repos.raw_file(args(body)?).await?),
2251 "raw_blobs" => reply(&repos.raw_blobs(args(body)?).await?),
2252 "visibility" => {
2253 let a: g1t_contracts::repos::VisibilityArgs = args(body)?;
2254 reply(&repos.registry.visibility(&a.paths).await?)
2255 }
2256 "storage" => reply(&repos.registry.storage().await?),
2257 "git_operations" => {
2258 let a: GitOperationsArgs = args(body)?;
2259 reply(&git_ops::totals(&repos.registry.db, &a.month, a.since.as_deref(), a.namespace.as_deref().map(str::to_lowercase).as_deref()).await?)
2260 }
2261 "all_ids" => {
2262 let a: AllIdsArgs = args(body)?;
2263 let limit = a.limit.clamp(1, 500);
2264 let ids = repos.registry.ids_after(a.after.as_deref(), limit).await?;
2265 let next = (ids.len() == limit as usize).then(|| ids.last().cloned()).flatten();
2266 reply(&IdPage { ids, next })
2267 }
2268 // The raw meters of the git store, for reconciling with Cloudflare
2269 // (meters.rs, scripts/ops/artifacts-usage.mjs).
2270 "artifacts_usage" => {
2271 let a: meters::UsageArgs = args(body)?;
2272 reply(&meters::usage(&repos.registry.db, &a).await?)
2273 }
2274 "operation_mapping" => reply(&meters::read_mapping(&repos.registry.db).await?),
2275 // Billing: the workspace each pull request's working copy is counted
2276 // for, so Cloudflare's own count of `pulls--<id>` shares out too.
2277 "pull_owners" => {
2278 #[derive(serde::Deserialize)]
2279 struct PullOwnersArgs {
2280 pulls: Vec<String>,
2281 }
2282 let a: PullOwnersArgs = args(body)?;
2283 let pulls: Vec<String> = a.pulls.into_iter().take(500).collect();
2284 reply(&serde_json::json!({ "owners": meters::pull_owners(&repos.registry.db, &pulls).await? }))
2285 }
2286 // Services only: which meters are operations, changed without a deploy.
2287 "set_operation_mapping" => {
2288 let row: meters::MappingRow = args(body)?;
2289 meters::set_mapping(&repos.registry.db, &row, &rfc3339(now_ms())).await?;
2290 reply(&meters::read_mapping(&repos.registry.db).await?)
2291 }
2292 // Backups (backups.rs): the runner's sweep claims queued ones, and
2293 // each sandbox, through the API, asks for its job and says how it went.
2294 "claim_backups" => {
2295 let a: g1t_contracts::backups::ClaimBackupsArgs = args(body)?;
2296 let blobs = backups::storage(&env);
2297 reply(&backups::claim(&repos.registry.db, blobs.as_ref(), &a, now_ms()).await?)
2298 }
2299 "backup_spec" => match backups::storage(&env) {
2300 Some(blobs) => {
2301 let a: g1t_contracts::backups::BackupJobArgs = args(body)?;
2302 let every = backups::Settings::from_env(&env).full_every;
2303 reply(&backups::spec(&repos.registry, &blobs, &repos.store, &a, every, now_ms()).await?)
2304 }
2305 None => reply(&backups_off::<bool>()),
2306 },
2307 "backup_complete" => match backups::storage(&env) {
2308 Some(blobs) => reply(&backups::complete(&repos.registry, &blobs, &args(body)?, now_ms()).await?),
2309 None => reply(&backups_off::<bool>()),
2310 },
2311 "backup_fail" => match backups::storage(&env) {
2312 Some(blobs) => reply(&backups::fail(&repos.registry.db, &blobs, &args(body)?).await?),
2313 None => reply(&backups_off::<bool>()),
2314 },
2315 // How the git store has been answering, for the status page.
2316 "store_health" => {
2317 let a: meters::HealthArgs = args(body)?;
2318 reply(&meters::health(&repos.registry.db, &a).await?)
2319 }
2320 // Where repositories may be kept, for a workspace's settings.
2321 "storage_options" => reply(&repos.storage_options()),
2322 // Services and operators only: how each namespace stands, and
2323 // moving a repository between them (namespaces.rs, moves.rs).
2324 "namespaces" => reply(&repos.standings().await?),
2325 "move_repository" => reply(&repos.move_repository(args(body)?).await?),
2326 "repository_moves" => {
2327 let a: moves::ListMovesArgs = args(body)?;
2328 reply(&repos.registry.moves(a.limit.unwrap_or(50)).await?)
2329 }
2330 _ => Response::error("Unknown method", 404),
2331 } }
2332 .await;
2333 // The git store is busy: said in words, with when to try again.
2334 let answered = match answered {
2335 Err(error) => match resilience::busy(&error.to_string()) {
2336 Some(busy) if OUTCOME_READS.contains(&method.as_str()) => {
2337 reply(&Outcome::<()>::fail(FailureCode::Conflict, busy.message().trim()))
2338 }
2339 Some(busy) => {
2340 let response = Response::error(busy.message(), 503)?;
2341 response.headers().set("retry-after", &busy.retry_after.to_string())?;
2342 Ok(response)
2343 }
2344 None => Err(error),
2345 },
2346 answered => answered,
2347 };
2348 flush_later(&env, &ctx);
2349 served.finish(answered)
2350}
2351
2352/// The nightly cron in wrangler.jsonc: tonight's backups are queued.
2353const BACKUP_CRON: &str = "53 2 * * *";
2354
2355/// The hourly sweep: deleted repositories whose time to be restored has
2356/// passed are purged. See lifecycle.rs. And, at [`BACKUP_CRON`], the
2357/// repositories whose refs moved are queued for a backup (backups.rs).
2358#[event(scheduled)]
2359async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
2360 let repos = match service(&env) {
2361 Ok(repos) => repos,
2362 Err(error) => {
2363 worker::console_error!("repos: the sweep could not start: {error}");
2364 return;
2365 }
2366 };
2367 if event.cron() == BACKUP_CRON {
2368 let Some(blobs) = backups::storage(&env) else { return };
2369 match backups::nightly(&repos.registry.db, &blobs, backups::Settings::from_env(&env), now_ms()).await {
2370 Ok(night) => worker::console_log!("repos: queued {} backups, removed {} of purged repositories", night.queued, night.pruned),
2371 Err(error) => worker::console_error!("repos: backups could not be queued: {error}"),
2372 }
2373 return;
2374 }
2375 match repos.purge_due(PurgeDueArgs::default()).await {
2376 Ok(0) => {}
2377 Ok(count) => worker::console_log!("repos: purged {count} deleted repositories"),
2378 Err(error) => worker::console_error!("repos: the purge sweep failed: {error}"),
2379 }
2380 // Pull requests' working copies whose time has come (forks.rs).
2381 match repos.retire_due().await {
2382 Ok(0) => {}
2383 Ok(count) => worker::console_log!("repos: removed {count} pull request working copies"),
2384 Err(error) => worker::console_error!("repos: the working copy sweep failed: {error}"),
2385 }
2386 // Repositories moving between namespaces, and old copies (moves.rs).
2387 match repos.run_moves().await {
2388 Ok(0) => {}
2389 Ok(count) => worker::console_log!("repos: moved {count} repositories between namespaces"),
2390 Err(error) => worker::console_error!("repos: the move sweep failed: {error}"),
2391 }
2392 meters::flush(&repos.registry.db).await;
2393}
2394
2395/// Events from the bus. A workspace's rename: its repositories move to the
2396/// workspace's current slug, asked of identity by id, so a repeated or late
2397/// delivery lands in the same place; their git store keys stay as they
2398/// were. A workspace's deletion: its repositories are deleted with it,
2399/// restored with it, or purged with it.
2400#[event(queue)]
2401async fn queue(batch: MessageBatch<Event>, env: Env, ctx: Context) -> Result<()> {
2402 let registry = Registry { db: env.d1("DB")? };
2403 let identity = env.service("IDENTITY")?;
2404 let handled = handle_events(&batch, &env, &registry, &identity).await;
2405 flush_later(&env, &ctx);
2406 handled
2407}
2408
2409async fn handle_events(batch: &MessageBatch<Event>, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
2410 for message in batch.messages()? {
2411 let event = message.body();
2412 // A pull request merged, closed or reopened: its working copy is
2413 // kept or let go (forks.rs).
2414 if let Some(change) = forks::pull_change(&event.kind) {
2415 let Some(pull_id) = forks::pull_id_of(&event.data) else {
2416 worker::console_error!("{} {} names no pull request", event.kind, event.id);
2417 continue;
2418 };
2419 let repos = service(env)?;
2420 match change {
2421 forks::PullChange::Settled => repos.pull_settled(&pull_id).await?,
2422 forks::PullChange::Reopened => repos.pull_reopened(&pull_id).await?,
2423 }
2424 continue;
2425 }
2426 // A workspace deleted, restored or purged: its repositories go with
2427 // it, come back with it, or are purged with it (lifecycle.rs).
2428 if event.kind == "workspace.deleting" {
2429 match serde_json::from_value::<WorkspaceDeleting>(event.data.clone()) {
2430 Ok(deleting) => service(env)?.delete_with_workspace(&deleting, &protected_workspaces(env)).await?,
2431 Err(_) => worker::console_error!("workspace.deleting {} could not be read", event.id),
2432 }
2433 continue;
2434 }
2435 if event.kind == "workspace.restored" {
2436 match serde_json::from_value::<WorkspaceRestored>(event.data.clone()) {
2437 Ok(restored) => service(env)?.restore_with_workspace(&restored).await?,
2438 Err(_) => worker::console_error!("workspace.restored {} could not be read", event.id),
2439 }
2440 continue;
2441 }
2442 if event.kind == "workspace.deleted" {
2443 match serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) {
2444 Ok(deleted) => service(env)?.purge_workspace(&deleted, &protected_workspaces(env)).await?,
2445 Err(_) => worker::console_error!("workspace.deleted {} could not be read", event.id),
2446 }
2447 continue;
2448 }
2449 if event.kind != "workspace.renamed" {
2450 continue;
2451 }
2452 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
2453 worker::console_error!("workspace.renamed {} could not be read", event.id);
2454 continue;
2455 };
2456 let names: HashMap<String, String> = g1t_kit::call(
2457 identity,
2458 "usernames",
2459 &g1t_contracts::identity::UsernamesArgs {
2460 ids: vec![renamed.workspace_id.clone()],
2461 },
2462 )
2463 .await?;
2464 let current = names
2465 .get(&renamed.workspace_id)
2466 .cloned()
2467 .unwrap_or_else(|| renamed.to.clone());
2468 let left = registry
2469 .rename_namespace(&renamed.stale_slugs(&current), &current)
2470 .await?;
2471 if left > 0 {
2472 worker::console_error!(
2473 "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
2474 renamed.from,
2475 renamed.to
2476 );
2477 }
2478 }
2479 Ok(())
2480}
2481
2482/// The workspaces whose repositories never go with a deletion, whatever is
2483/// published: `PROTECTED_WORKSPACES` if set here, and Flagon's always.
2484fn protected_workspaces(env: &Env) -> Vec<String> {
2485 let configured = env.var("PROTECTED_WORKSPACES").ok().map(|v| v.to_string());
2486 g1t_contracts::identity::protected_names(configured.as_deref())
2487}
2488
2489/// The repository a push to a path that does not exist yet creates: private,
2490/// so nothing pushed by mistake is published. An owner makes it public on
2491/// purpose (`POST /repos/{owner}/{repo}/visibility`).
2492fn push_to_create(owner: &User, path: &RepoPath) -> CreateArgs {
2493 CreateArgs {
2494 owner: owner.clone(),
2495 namespace: path.namespace.clone(),
2496 name: path.name.clone(),
2497 description: None,
2498 is_private: true,
2499 import_url: None,
2500 import_token: None,
2501 }
2502}
2503
2504#[cfg(test)]
2505mod push_to_create_tests {
2506 use super::*;
2507
2508 #[test]
2509 fn a_pushed_repository_starts_private() {
2510 let owner: User = serde_json::from_value(serde_json::json!({ "id": "usr_1", "username": "ada" })).unwrap();
2511 let args = push_to_create(&owner, &RepoPath { namespace: "acme".into(), name: "site".into() });
2512 assert!(args.is_private);
2513 assert_eq!((args.namespace.as_str(), args.name.as_str()), ("acme", "site"));
2514 }
2515}