g1t/services/repos/src/lib.rs

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