g1t/services/repos/src/lib.rs

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