g1t/services/repos/src/lib.rs

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