g1t/services/repos/src/lib.rs

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