g1t/services/repos/src/lib.rs

2,214 lines92,594 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 let (entries, complete) = last_commits::last_commits(&git, &head.hash, &a.tree_path).await?;
774 let found = g1t_contracts::repos::LastCommits { entries, complete };
775 if let Ok(mut response) = worker::Response::from_json(&found) {
776 let _ = response.headers_mut().set("cache-control", "max-age=604800");
777 let _ = cache.put(key.as_str(), response).await;
778 }
779 Ok(Outcome::Ok(found))
780 }
781
782 /// The repository's tags, newest commit first, at most 100.
783 async fn tags(&self, a: g1t_contracts::repos::TagsArgs) -> Result<Outcome<Vec<g1t_contracts::repos::Tag>>> {
784 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
785 return Ok(not_found());
786 };
787 let git = self.store.open(&store_key(&repo)).await?;
788 let access = git.access(Scope::Read).await?;
789 let named: Vec<(String, String)> = refs::heads_and_tags(refs::all(&access).await?)
790 .into_iter()
791 .filter_map(|(name, hash)| name.strip_prefix("refs/tags/").map(|tag| (tag.to_owned(), hash)))
792 .collect();
793 let read = self.read_git(&repo).await?;
794 let commits = futures_util::future::join_all(named.iter().take(MAX_TAGS_READ).map(|(_, hash)| read.log(hash, 1))).await;
795 let mut tags: Vec<g1t_contracts::repos::Tag> = named
796 .into_iter()
797 .zip(commits.into_iter().map(|found| found.ok().and_then(|list| list.into_iter().next())).chain(std::iter::repeat(None)))
798 .map(|((name, _), commit)| g1t_contracts::repos::Tag { name, commit })
799 .collect();
800 tags.sort_by(|a, b| {
801 let at = |tag: &g1t_contracts::repos::Tag| tag.commit.as_ref().map(|c| c.authored_at.clone()).unwrap_or_default();
802 at(b).cmp(&at(a)).then_with(|| b.name.cmp(&a.name))
803 });
804 tags.truncate(MAX_TAGS_READ);
805 Ok(Outcome::Ok(tags))
806 }
807
808 /// The repository's branches, default branch first.
809 async fn branches(&self, a: BranchesArgs) -> Result<Outcome<Vec<Branch>>> {
810 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
811 return Ok(not_found());
812 };
813 let mut branches = self.read_git(&repo).await?.branches().await?;
814 branches.sort_by_key(|branch| branch.name != repo.default_branch);
815 Ok(Outcome::Ok(branches))
816 }
817
818 /// Whether a pull request's source lacks commits that the branch it
819 /// would merge into has.
820 async fn behind(&self, a: BehindArgs) -> Result<bool> {
821 let Some(source) = self.registry.by_id(&a.source_id).await? else {
822 return Ok(false);
823 };
824 let target = match &source.fork_of {
825 Some(id) => self.registry.by_id(id).await?,
826 None => Some(source.clone()),
827 };
828 let Some(target) = target else {
829 return Ok(false);
830 };
831 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
832 let target_head = self
833 .read_git(&target)
834 .await?
835 .log(&target.default_branch, 1)
836 .await?
837 .into_iter()
838 .next()
839 .map(|commit| commit.hash);
840 let Some(target_head) = target_head else {
841 return Ok(false);
842 };
843 let source_git = self.read_git(&source).await?;
844 let history = source_git.log(&branch, MAX_ANCESTRY).await?;
845 if history.is_empty() {
846 return Ok(false);
847 }
848 Ok(!descends_from(&source_git, &history, &target_head).await?)
849 }
850
851 /// The files a pull request's source and the default branch it would
852 /// merge into each changed since they last agreed. Where the two lists
853 /// share no file, the merge cannot conflict; where they do, it may.
854 async fn divergence(&self, a: BehindArgs) -> Result<Option<Divergence>> {
855 let Some(source) = self.registry.by_id(&a.source_id).await? else {
856 return Ok(None);
857 };
858 let target = match &source.fork_of {
859 Some(id) => self.registry.by_id(id).await?,
860 None => Some(source.clone()),
861 };
862 let Some(target) = target else {
863 return Ok(None);
864 };
865 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
866 let source_git = self.read_git(&source).await?;
867 let target_git = self.read_git(&target).await?;
868 // The target's side is the same for every pull request into it, and
869 // worked out once per head (coalesce.rs).
870 let (history, side) = futures_util::future::try_join(
871 source_git.log(&branch, MAX_ANCESTRY),
872 self.target_side(&target, &target_git),
873 )
874 .await?;
875 let target_history = &side.history;
876 let (Some(head), Some(base)) = (history.first(), target_history.first()) else {
877 return Ok(None);
878 };
879 let behind = !descends_from(&source_git, &history, &base.hash).await?;
880 let merge_base = nearest_ancestor_in(&source_git, &history, &side.shared).await?;
881 let mut divergence = Divergence {
882 head: head.hash.clone(),
883 base: base.hash.clone(),
884 merge_base: merge_base.clone(),
885 behind,
886 ..Divergence::default()
887 };
888 let merge_base_tree = match &merge_base {
889 Some(hash) => target_history
890 .iter()
891 .find(|commit| commit.hash == *hash)
892 .map(|commit| commit.tree_hash.clone()),
893 None => None,
894 };
895 let Some(merge_base_tree) = merge_base_tree else {
896 // No common history to compare from: say nothing is known.
897 divergence.truncated = true;
898 return Ok(Some(divergence));
899 };
900 let (ours, truncated_ours) =
901 diff::changed_paths(&source_git, Some(&merge_base_tree), &head.tree_hash).await?;
902 divergence.ours = ours;
903 divergence.truncated = truncated_ours;
904 if behind {
905 let now = now_ms();
906 let key = (target.id.clone(), merge_base_tree.clone(), base.tree_hash.clone());
907 let (theirs, truncated_theirs) = match THEIRS.with(|memo| memo.borrow().get(&key, now)) {
908 Some(kept) => kept,
909 None => {
910 let found = diff::changed_paths(&target_git, Some(&merge_base_tree), &base.tree_hash).await?;
911 THEIRS.with(|memo| memo.borrow_mut().put(key, found.clone(), now));
912 found
913 }
914 };
915 divergence.theirs = theirs;
916 divergence.truncated |= truncated_theirs;
917 }
918 Ok(Some(divergence))
919 }
920
921 /// A target branch's history from its head, worked out once per head
922 /// for every pull request asking about it (coalesce.rs). The head is
923 /// read under the refs version; the history by its hash, which the
924 /// object cache keeps for good.
925 async fn target_side<R: GitRepo>(&self, target: &Repo, git: &R) -> Result<Rc<coalesce::TargetSide>> {
926 let now = now_ms();
927 let key = refs_cache::usable(registry::refs_state(&target.id), now)
928 .map(|version| (target.id.clone(), target.default_branch.clone(), version));
929 if let Some(key) = &key
930 && let Some(side) = TARGETS.with(|memo| memo.borrow().get(key, now))
931 {
932 return Ok(side);
933 }
934 let history = match git.log(&target.default_branch, 1).await?.first() {
935 Some(head) => git.log(&head.hash, MAX_ANCESTRY).await?,
936 None => Vec::new(),
937 };
938 let side = coalesce::TargetSide::new(history);
939 if let Some(key) = key {
940 TARGETS.with(|memo| memo.borrow_mut().put(key, side.clone(), now));
941 }
942 Ok(side)
943 }
944
945 async fn head(&self, a: HeadArgs) -> Result<Option<String>> {
946 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
947 return Ok(None);
948 };
949 let branch = if a.branch.is_empty() { &repo.default_branch } else { &a.branch };
950 let git = self.read_git(&repo).await?;
951 Ok(git
952 .log(branch, 1)
953 .await?
954 .into_iter()
955 .next()
956 .map(|commit| commit.hash))
957 }
958
959 async fn delete_branch(&self, a: DeleteBranchArgs) -> Result<Outcome<bool>> {
960 if !a.branch.starts_with(G1T_BRANCH_PREFIX) {
961 return Ok(Outcome::fail(
962 FailureCode::Forbidden,
963 "Only branches g1t made for itself can be deleted this way.",
964 ));
965 }
966 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
967 return Ok(not_found());
968 };
969 self.live(&repo).await?;
970 let git = self.store.open(&store_key(&repo)).await?;
971 let Some(old) = git
972 .branches()
973 .await?
974 .into_iter()
975 .find(|branch| branch.name == a.branch)
976 .map(|branch| branch.hash)
977 else {
978 return Ok(Outcome::Ok(false));
979 };
980 let access = git.access(Scope::Write).await?;
981 let deleted = land::delete_ref(&access, &a.branch, &old).await?;
982 self.refs_moved(&repo.id).await;
983 if let Err(reason) = deleted {
984 return Ok(Outcome::fail(
985 FailureCode::Conflict,
986 format!("{} could not be deleted: {reason}", a.branch),
987 ));
988 }
989 Ok(Outcome::Ok(true))
990 }
991
992 async fn fork_for_pull(&self, a: ForkArgs) -> Result<Outcome<Repo>> {
993 let viewer = Some(a.actor.clone());
994 let Some(source) = self
995 .registry
996 .by_id(&a.source_id)
997 .await?
998 .filter(|repo| can_read(repo, &viewer))
999 else {
1000 return Ok(not_found());
1001 };
1002 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1003 return Ok(Outcome::fail(code, message));
1004 }
1005 let now = now_ms();
1006 let fork = Repo {
1007 id: new_id("rep", now),
1008 namespace: PULLS_NAMESPACE.to_owned(),
1009 name: a.pull_id.clone(),
1010 description: None,
1011 // A fork is exactly as visible as the repo it came from.
1012 is_private: source.is_private,
1013 owner_id: a.actor.id.clone(),
1014 default_branch: source.default_branch.clone(),
1015 fork_of: Some(source.id.clone()),
1016 protected: false,
1017 created_at: rfc3339(now),
1018 topics: Vec::new(),
1019 website: None,
1020 archived_at: None,
1021 };
1022 // Artifacts forks within a namespace: the copy goes where its
1023 // repository is.
1024 let (namespace, _) = store::locate(&store_key(&source));
1025 self.registry
1026 .claim_store_key(&fork, Some(&namespace), &self.store.default_namespace())
1027 .await?;
1028 self.store
1029 .open(&store_key(&source))
1030 .await?
1031 .fork(&store_key(&fork))
1032 .await?;
1033 self.registry.insert(&fork).await?;
1034 self.publish(NewEvent {
1035 kind: "repo.forked",
1036 source: SOURCE,
1037 repo_id: Some(source.id.clone()),
1038 actor: Some(a.actor.id),
1039 data: RepoForked {
1040 repo_id: fork.id.clone(),
1041 source_repo_id: source.id,
1042 pull_id: a.pull_id,
1043 },
1044 })
1045 .await?;
1046 Ok(Outcome::Ok(fork))
1047 }
1048
1049 async fn git_access(&self, a: GitAccessArgs) -> Result<Outcome<GitAccess>> {
1050 let found = self.registry.by_path(&a.path).await?;
1051 Ok(match self.authorize_git(&a.path, &a.viewer, a.service, found).await? {
1052 Outcome::Ok(repo) => {
1053 self.live(&repo).await?;
1054 let write = a.service == GitService::ReceivePack;
1055 if write {
1056 // A push with this credential would not pass through
1057 // here, so nothing that lists the refs is kept until it
1058 // has expired (see refs_cache.rs).
1059 let until = now_ms() + store::CREDENTIAL_LIFE_MS + 60_000;
1060 if let Err(error) = self.registry.refs_open(&repo.id, until).await {
1061 // Before the column exists nothing is kept anyway.
1062 if registry::refs_state(&repo.id).is_some() {
1063 return Err(error);
1064 }
1065 }
1066 }
1067 let scope = if write { Scope::Write } else { Scope::Read };
1068 Outcome::Ok(self.store.handout(&store_key(&repo), scope).await?)
1069 }
1070 Outcome::Fail(failure) => Outcome::Fail(failure),
1071 })
1072 }
1073
1074 /// The repository at `path` (`found`, as just read), if the viewer may
1075 /// use `service` on it: fetch from it, or push to it. A push to a path
1076 /// with nothing there makes the repository, in a workspace the pusher
1077 /// belongs to.
1078 async fn authorize_git(
1079 &self,
1080 path: &RepoPath,
1081 viewer: &Viewer,
1082 service: GitService,
1083 found: Option<Repo>,
1084 ) -> Result<Outcome<Repo>> {
1085 let mut a = GitAccessArgs {
1086 path: path.clone(),
1087 viewer: viewer.clone(),
1088 service,
1089 };
1090 let write = a.service == GitService::ReceivePack;
1091 // An access token: pushing needs code:write, reading a private
1092 // repository code:read. A public repository reads as it would for
1093 // anyone. Which repositories a token reaches is its owner's, checked
1094 // below as for anyone.
1095 if let Some(access) = a.viewer.as_ref().and_then(|user| user.token.as_deref()).cloned() {
1096 let public = found.as_ref().is_some_and(|repo| !repo.is_private);
1097 let decision = g1t_contracts::scopes::decide_git(&access, write, public);
1098 if !decision.allowed {
1099 return Ok(Outcome::fail(
1100 FailureCode::Forbidden,
1101 format!("{}\n", decision.reason.unwrap_or_default()),
1102 ));
1103 }
1104 if !write && !access.allows(g1t_contracts::scopes::Scope::CodeRead) {
1105 a.viewer = None;
1106 }
1107 }
1108
1109 // Anonymous callers are asked to authenticate whether or not the repo
1110 // exists, so private repos cannot be told apart from missing ones.
1111 let denied = || match &a.viewer {
1112 Some(_) => not_found(),
1113 None => Outcome::fail(FailureCode::Unauthenticated, "Authentication required."),
1114 };
1115 // An agent's token works through the API only: its sandbox has its
1116 // own way to push, to its own pull request.
1117 if a.viewer.as_ref().is_some_and(|user| user.kind == PrincipalKind::Agent) {
1118 return Ok(Outcome::fail(
1119 FailureCode::Forbidden,
1120 "A g1t agent's token cannot be used with git.",
1121 ));
1122 }
1123 if let (true, Some(user)) = (write, &a.viewer)
1124 && !user.verified
1125 {
1126 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1127 }
1128 let repo = match found {
1129 Some(repo) => {
1130 let allowed = if write {
1131 can_write(&repo, &a.viewer)
1132 } else {
1133 self.may_read(&repo, &a.viewer).await?
1134 };
1135 if !allowed {
1136 return Ok(denied());
1137 }
1138 // An archived repository, or a pull request's copy of one,
1139 // is read-only.
1140 if write {
1141 let archived = match &repo.fork_of {
1142 Some(source) => self.registry.by_id(source).await?,
1143 None => Some(repo.clone()),
1144 };
1145 match archived {
1146 Some(source) => {
1147 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1148 return Ok(Outcome::fail(code, format!("{message}\n")));
1149 }
1150 }
1151 // The repository it was copied from is deleted.
1152 None => return Ok(denied()),
1153 }
1154 }
1155 repo
1156 }
1157 None => {
1158 // Push to create, in a workspace the pusher belongs to.
1159 let owner = a
1160 .viewer
1161 .as_ref()
1162 .filter(|user| write && user.is_member(&a.path.namespace.to_lowercase()));
1163 let Some(owner) = owner else {
1164 return Ok(denied());
1165 };
1166 let created = self.create(push_to_create(owner, &a.path)).await?;
1167 match created {
1168 Outcome::Ok(repo) => repo,
1169 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
1170 }
1171 }
1172 };
1173 Ok(Outcome::Ok(repo))
1174 }
1175
1176 async fn land(&self, a: LandArgs) -> Result<Outcome<Landed>> {
1177 let actor: Viewer = Some(a.actor.clone());
1178 let Some(source) = self.registry.by_id(&a.source_id).await? else {
1179 return Ok(not_found());
1180 };
1181 // A fork lands on the repository it came from; a branch on its own.
1182 let target = match &source.fork_of {
1183 Some(id) => self.registry.by_id(id).await?,
1184 None => Some(source.clone()),
1185 };
1186 let Some(target) = target.filter(|repo| can_read(repo, &actor)) else {
1187 return Ok(not_found());
1188 };
1189 if !registry::can(&target, &actor, Capability::Merge) {
1190 return Ok(Outcome::fail(
1191 FailureCode::Forbidden,
1192 access::needs(Capability::Merge, &format!("{}/{}", target.namespace, target.name)),
1193 ));
1194 }
1195 if !a.actor.verified {
1196 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1197 }
1198 if let Some((code, message)) = lifecycle::archived_refusal(&target) {
1199 return Ok(Outcome::fail(code, message));
1200 }
1201
1202 let branch = &target.default_branch;
1203 let from_fork = source.id != target.id;
1204 let source_branch = match a.branch {
1205 Some(name) if !from_fork && name == *branch => {
1206 return Ok(Outcome::fail(
1207 FailureCode::Invalid,
1208 format!("{branch} cannot be merged into itself."),
1209 ));
1210 }
1211 Some(name) => name,
1212 None if from_fork => branch.clone(),
1213 None => {
1214 return Ok(Outcome::fail(
1215 FailureCode::Invalid,
1216 "Say which branch to merge.",
1217 ));
1218 }
1219 };
1220
1221 self.live(&source).await?;
1222 let source_git = self.store.open(&store_key(&source)).await?;
1223 let target_git = self.store.open(&store_key(&target)).await?;
1224 let history = source_git.log(&source_branch, MAX_ANCESTRY).await?;
1225 let Some(new) = history.first().map(|commit| commit.hash.clone()) else {
1226 return Ok(Outcome::fail(
1227 FailureCode::Conflict,
1228 "This pull request has no commits to merge.",
1229 ));
1230 };
1231 let old = target_git
1232 .log(branch, 1)
1233 .await?
1234 .into_iter()
1235 .next()
1236 .map(|commit| commit.hash);
1237
1238 if old.as_deref() == Some(new.as_str()) {
1239 return Ok(Outcome::Ok(Landed {
1240 commit: new,
1241 previous: None,
1242 }));
1243 }
1244 // Moving the branch to a commit that does not descend from its
1245 // current head would discard whatever landed in between.
1246 if let Some(old) = &old
1247 && !descends_from(&source_git, &history, old).await?
1248 {
1249 let remedy = if from_fork {
1250 format!("Pull {branch} into the pull request's fork, push, and merge again.")
1251 } else {
1252 format!("Merge {branch} into {source_branch}, push, and merge again.")
1253 };
1254 return Ok(Outcome::fail(
1255 FailureCode::Conflict,
1256 format!("{branch} has moved since this pull request was opened. {remedy}"),
1257 ));
1258 }
1259
1260 // For a branch the objects are already in the target; sending them
1261 // again is harmless and keeps one way of moving a ref.
1262 let source_access = source_git.access(Scope::Read).await?;
1263 let target_access = target_git.access(Scope::Write).await?;
1264 let pushed =
1265 land::fast_forward(&source_access, &target_access, branch, old.as_deref(), &new)
1266 .await?;
1267 self.refs_moved(&target.id).await;
1268 if let Err(reason) = pushed {
1269 // Most often another pull request landed between the check and the push.
1270 return Ok(Outcome::fail(
1271 FailureCode::Conflict,
1272 format!("{branch} could not be updated: {reason}"),
1273 ));
1274 }
1275 self.publish_push(
1276 &target,
1277 &format!("refs/heads/{branch}"),
1278 old.as_deref(),
1279 &new,
1280 Some(a.actor.id),
1281 )
1282 .await?;
1283 Ok(Outcome::Ok(Landed {
1284 commit: new,
1285 previous: old,
1286 }))
1287 }
1288
1289 async fn compare(&self, a: CompareArgs) -> Result<Outcome<Comparison>> {
1290 let Some(repo) = self
1291 .visible(self.registry.by_id(&a.repo_id).await?, &a.viewer)
1292 .await?
1293 else {
1294 return Ok(not_found());
1295 };
1296 let git = self.read_git(&repo).await?;
1297 let head_ref = a.head.as_deref().unwrap_or(&repo.default_branch);
1298 // The head's history is only searched when the base is worked out
1299 // from another branch.
1300 let depth = if a.base.is_some() || is_commit_hash(head_ref) { 1 } else { MAX_ANCESTRY };
1301 let history = git.log(head_ref, depth).await?;
1302 let Some(head) = history.first() else {
1303 return Ok(Outcome::fail(
1304 FailureCode::Conflict,
1305 "There are no commits to compare.",
1306 ));
1307 };
1308
1309 // Where the head's history meets the default branch of `against`.
1310 let shared_with = async |against: &Repo| -> Result<Option<String>> {
1311 let against_git = self.read_git(against).await?;
1312 let shared: HashSet<String> = against_git
1313 .log(&against.default_branch, MAX_ANCESTRY)
1314 .await?
1315 .into_iter()
1316 .map(|commit| commit.hash)
1317 .collect();
1318 nearest_ancestor_in(&git, &history, &shared).await
1319 };
1320 let base = match (a.base, &repo.fork_of) {
1321 (Some(base), _) => Some(base),
1322 // A fork is compared with the last commit it shares with the
1323 // repository it came from.
1324 (None, Some(target_id)) => match self.registry.by_id(target_id).await? {
1325 Some(target) => shared_with(&target).await?,
1326 None => None,
1327 },
1328 // A branch, with the point where it left the default branch.
1329 // A single commit, with its first parent.
1330 (None, None) if is_commit_hash(head_ref) => head.parents.first().cloned(),
1331 (None, None) if head_ref != repo.default_branch => shared_with(&repo).await?,
1332 (None, None) => head.parents.first().cloned(),
1333 };
1334 let base_tree = match &base {
1335 Some(base) => git
1336 .log(base, 1)
1337 .await?
1338 .into_iter()
1339 .next()
1340 .map(|commit| commit.tree_hash),
1341 None => None,
1342 };
1343 let (files, truncated) =
1344 diff::compare_trees(&git, base_tree.as_deref(), &head.tree_hash).await?;
1345 Ok(Outcome::Ok(Comparison {
1346 base,
1347 head: head.hash.clone(),
1348 files,
1349 truncated,
1350 }))
1351 }
1352
1353 /// Reports that `git_ref` of `repo` (a full ref) now points to `after`.
1354 async fn publish_push(
1355 &self,
1356 repo: &Repo,
1357 git_ref: &str,
1358 before: Option<&str>,
1359 after: &str,
1360 actor: Option<String>,
1361 ) -> Result<()> {
1362 self.publish_git_push(repo, git_ref, before, after, actor, false).await
1363 }
1364
1365 /// `publish_push`, saying whether the push reached the store without
1366 /// being scanned for secrets first.
1367 async fn publish_git_push(
1368 &self,
1369 repo: &Repo,
1370 git_ref: &str,
1371 before: Option<&str>,
1372 after: &str,
1373 actor: Option<String>,
1374 unscanned: bool,
1375 ) -> Result<()> {
1376 self.publish(NewEvent {
1377 kind: "git.push",
1378 source: SOURCE,
1379 repo_id: Some(repo.id.clone()),
1380 actor,
1381 data: GitPush {
1382 repo_id: repo.id.clone(),
1383 git_ref: git_ref.to_owned(),
1384 before: before.map(str::to_owned),
1385 after: after.to_owned(),
1386 default_branch: git_ref.strip_prefix("refs/heads/")
1387 == Some(repo.default_branch.as_str()),
1388 unscanned,
1389 },
1390 })
1391 .await
1392 }
1393
1394 /// Git over HTTPS. Only what decides the answer happens before it:
1395 /// the repository, who is asking and whether they may, the free
1396 /// workspace limits, push protection, and the store's own answer. The
1397 /// audit entry and what a push changed are recorded once git has its
1398 /// answer. Each answer says how long its steps took (`Server-Timing`).
1399 async fn git_http(&self, request: Request, env: &Env, ctx: &Context) -> Result<Response> {
1400 let mut timing = git_http::Timing::start();
1401 let Some(git) = git_http::parse(&request.url()?) else {
1402 return Response::error("Not found", 404);
1403 };
1404 let response = match self.answer_git(request, &git, env, ctx, &mut timing).await {
1405 Ok(response) => response,
1406 // The git store is busy: git hears when to try again.
1407 Err(error) => match resilience::busy(&error.to_string()) {
1408 Some(busy) => git_http::busy_response(busy)?,
1409 None => return Err(error),
1410 },
1411 };
1412 timing.apply(response)
1413 }
1414
1415 async fn answer_git(
1416 &self,
1417 request: Request,
1418 git: &git_http::GitRequest,
1419 env: &Env,
1420 ctx: &Context,
1421 timing: &mut git_http::Timing,
1422 ) -> Result<Response> {
1423 let write = git.service == GitService::ReceivePack;
1424 let get = request.method() == Method::Get;
1425 let identity = env.service("IDENTITY")?;
1426 // The repository and the caller's credentials, at once. A fetch may
1427 // go by the row as read a moment ago, for the same clone's next
1428 // request; a push always reads it. Anonymous callers cost nothing.
1429 let lookup = async {
1430 if write {
1431 self.registry.by_path(&git.path).await
1432 } else {
1433 self.registry.by_path_recent(&git.path).await
1434 }
1435 };
1436 let (found, viewer) =
1437 futures_util::future::join(lookup, git_http::viewer(&request, &identity)).await;
1438 let found = found?;
1439 timing.mark("repo");
1440 if found.is_none() {
1441 // A workspace that was renamed: git follows a redirect when it
1442 // first asks for refs, and uses the new address from then on.
1443 // A repository transferred to another workspace: the same, to
1444 // its new path. Fetches and pushes both follow either.
1445 let url = request.url()?;
1446 let (renamed, moved) = futures_util::future::join(
1447 git_http::renamed(&url, &identity),
1448 self.registry.resolve_moved(&git.path),
1449 )
1450 .await;
1451 timing.mark("moved");
1452 if let Some(location) = renamed? {
1453 return git_http::moved(&location, get);
1454 }
1455 if let Some(now) = moved?
1456 && let Some(location) = git_http::transferred(&url, &now)
1457 {
1458 return git_http::moved(&location, get);
1459 }
1460 }
1461 let viewer = viewer?;
1462 // A run credential is checked against its grants, then acts as the
1463 // person it works for. See run_access.rs.
1464 let (request, viewer, audit) = match self.admit_git(request, git, viewer, found.as_ref()).await? {
1465 run_access::Admitted::Go { request, viewer, entry } => (request, viewer, entry),
1466 run_access::Admitted::Refused(response) => return Ok(response),
1467 };
1468 let mut after = AfterGit {
1469 audit,
1470 status: 0,
1471 message: None,
1472 push: None,
1473 };
1474 let repo = match self.authorize_git(&git.path, &viewer, git.service, found).await? {
1475 Outcome::Ok(repo) => repo,
1476 refused => {
1477 let response = git_http::refuse(refused)?;
1478 after.ended(response.status_code(), None);
1479 after.spawn(env, ctx);
1480 return Ok(response);
1481 }
1482 };
1483 // A pull request's working copy removed after it closed is made
1484 // again before git uses it (forks.rs).
1485 self.live(&repo).await?;
1486 timing.mark("access");
1487 // A protected default branch takes changes only from a merged pull
1488 // request, which lands without going through here.
1489 let protected = (repo.protected && repo.fork_of.is_none()).then(|| repo.default_branch.clone());
1490 // Clones check out the default branch g1t keeps, which can have
1491 // changed since the store made the repository.
1492 let default_branch = repo.fork_of.is_none().then(|| repo.default_branch.clone());
1493 let key = store_key(&repo);
1494 let scope = if write { Scope::Write } else { Scope::Read };
1495 let mut request = request;
1496 let protocol = refs_cache::protocol(request.headers().get("git-protocol")?.as_deref());
1497 // A fetch's POST is read here, to tell an `ls-refs` from a fetch of
1498 // objects; the store would have it read in full anyway.
1499 let body = if !write && !get { Some(request.bytes().await?) } else { None };
1500 // What it asks the store, for the meters (meters.rs).
1501 let call = git_ops::classify(git.service, git.endpoint, get, body.as_deref());
1502 // An answer that lists refs may have been kept: see refs_cache.rs.
1503 let kept_key = refs_cache::kind(git, get, protocol, body.as_deref())
1504 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1505 .map(|(kind, version)| {
1506 refs_cache::Key::new(&repo.id, version, default_branch.as_deref(), protocol, &kind)
1507 });
1508 // A kept answer and the free workspace limits, with a kept
1509 // credential looked up alongside. A kept answer goes back without
1510 // waiting for the credential, which it does not need.
1511 let ((answer, limited), kept_access) = {
1512 let shared = self.shared.as_deref();
1513 let answer_and_limits = std::pin::pin!(futures_util::future::join(
1514 async {
1515 match &kept_key {
1516 Some(kept_key) => refs_cache::get(shared, kept_key).await,
1517 None => None,
1518 }
1519 },
1520 self.git_limits(call, git, &repo, env),
1521 ));
1522 let kept_access = std::pin::pin!(self.store.kept_access(&key, scope));
1523 match futures_util::future::select(answer_and_limits, kept_access).await {
1524 futures_util::future::Either::Left((first, kept_access)) => {
1525 let answered = first.0.is_some() || matches!(first.1, Ok(Some(_)) | Err(_));
1526 (first, if answered { None } else { kept_access.await })
1527 }
1528 futures_util::future::Either::Right((kept_access, first)) => (first.await, kept_access),
1529 }
1530 };
1531 timing.mark("kept");
1532 if let Some((response, status, message)) = limited? {
1533 after.ended(status, Some(message.to_owned()));
1534 after.spawn(env, ctx);
1535 return Ok(response);
1536 }
1537 if let (Some((entry, found)), Some(kept_key)) = (answer, &kept_key) {
1538 timing.note("refs", found.as_str());
1539 if found == refs_cache::Found::Shared {
1540 let (kept_key, entry) = (kept_key.clone(), entry.clone());
1541 ctx.wait_until(async move { refs_cache::keep_in_colo(&kept_key, &entry).await });
1542 }
1543 // Never reached the store: never an operation.
1544 meters::record(call.cached_meter(), &key, 0, entry.body.len() as u64);
1545 after.ended(200, None);
1546 after.spawn(env, ctx);
1547 return entry.response();
1548 }
1549 if kept_key.is_some() {
1550 timing.note("refs", "miss");
1551 }
1552 // The store's credential: one made a moment ago, here or in another
1553 // isolate (see store.rs), or a new one.
1554 let access = match kept_access {
1555 Some((access, from)) => {
1556 timing.note("cred", from.as_str());
1557 access
1558 }
1559 None => {
1560 let access = self.store.mint_access(&key, scope).await?;
1561 timing.mark("mint");
1562 timing.note("cred", "mint");
1563 access
1564 }
1565 };
1566 // Should the store turn a kept credential down, a fetch's first
1567 // request is tried again with a new one; the requests after it then
1568 // have that one too.
1569 let again = if get { Some(request.clone()?) } else { None };
1570 // Push protection: a push that adds a secret is refused. See secret_scan.rs.
1571 let scan = async |body: &[u8]| self.protect(&repo, viewer.as_ref(), body).await;
1572 // What a push may bring (pack_limits.rs): the repository's size is
1573 // its own and its pull requests' working copies'.
1574 let limits = if write && !get {
1575 git_http::PushLimits {
1576 held: self.held(&repo).await,
1577 repo_limit: self.repo_limit,
1578 large: self.large_pushes,
1579 ..git_http::PushLimits::default()
1580 }
1581 } else {
1582 git_http::PushLimits::default()
1583 };
1584 let mut outcome = git_http::forward(
1585 request,
1586 body,
1587 git,
1588 &access,
1589 protected.as_deref(),
1590 default_branch.as_deref(),
1591 limits,
1592 scan,
1593 )
1594 .await?;
1595 let turned_down = matches!(
1596 &outcome,
1597 git_http::Push::Forwarded(forwarded) if matches!(forwarded.response.status_code(), 401 | 403)
1598 );
1599 if turned_down {
1600 self.store.forget_access(&key).await;
1601 if let Some(again) = again {
1602 let access = self.store.mint_access(&key, scope).await?;
1603 let nothing = async |_: &[u8]| Ok(None);
1604 outcome = git_http::forward(
1605 again,
1606 None,
1607 git,
1608 &access,
1609 protected.as_deref(),
1610 default_branch.as_deref(),
1611 git_http::PushLimits::default(),
1612 nothing,
1613 )
1614 .await?;
1615 }
1616 }
1617 let forwarded =
1618 match outcome {
1619 git_http::Push::Forwarded(forwarded) => forwarded,
1620 git_http::Push::Refused(response) => {
1621 after.ended(403, Some("The push would change a protected branch.".to_owned()));
1622 after.spawn(env, ctx);
1623 return Ok(response);
1624 }
1625 git_http::Push::Blocked(response) => {
1626 after.ended(403, Some("The push adds a secret.".to_owned()));
1627 after.spawn(env, ctx);
1628 return Ok(response);
1629 }
1630 git_http::Push::Declined(response, reason) => {
1631 after.ended(403, Some(format!("The push was declined: {reason}.")));
1632 after.spawn(env, ctx);
1633 return Ok(response);
1634 }
1635 };
1636 if forwarded.from_store {
1637 let received = forwarded
1638 .response
1639 .headers()
1640 .get("content-length")?
1641 .and_then(|length| length.parse().ok())
1642 .unwrap_or(0);
1643 meters::record(call.meter(), &key, forwarded.sent, received);
1644 }
1645 timing.mark("store");
1646 let mut response = forwarded.response;
1647 let status = response.status_code();
1648 if write && !get {
1649 // A push: the store has moved its refs once it has answered in
1650 // full, so the answer is read before the change is recorded, and
1651 // only then goes back. Whoever fetches after it sees the push.
1652 let headers = response.headers().clone();
1653 headers.delete("content-length")?;
1654 let report = response.bytes().await?;
1655 self.refs_moved(&repo.id).await;
1656 timing.mark("refs");
1657 response = Response::from_bytes(report)?.with_headers(headers).with_status(status);
1658 } else if let (Some(kept_key), 200) = (&kept_key, status) {
1659 // A miss: this answer is kept for the next to ask.
1660 let headers = response.headers().clone();
1661 headers.delete("content-length")?;
1662 let body = response.bytes().await?;
1663 if let Some(content_type) = headers.get("content-type")? {
1664 let entry = refs_cache::Entry { content_type, body: body.clone() };
1665 if entry.keepable() {
1666 let shared = self.shared.clone();
1667 let kept_key = kept_key.clone();
1668 ctx.wait_until(async move { refs_cache::keep(shared.as_deref(), &kept_key, &entry).await });
1669 }
1670 }
1671 response = Response::from_bytes(body)?.with_headers(headers).with_status(status);
1672 }
1673 after.ended(status, None);
1674 if status == 200 && (forwarded.pack_bytes > 0 || !forwarded.pushed.is_empty()) {
1675 after.push = Some(PushDone {
1676 repo,
1677 pushed: forwarded.pushed,
1678 pack_bytes: forwarded.pack_bytes,
1679 actor: viewer.map(|user: User| user.id),
1680 unscanned: forwarded.unscanned,
1681 });
1682 }
1683 after.spawn(env, ctx);
1684 Ok(response)
1685 }
1686
1687 /// The answer for a request a free workspace's limits stop, or a push
1688 /// to a full repository, with its status and reason for the audit log;
1689 /// `None` to go on.
1690 ///
1691 /// A clone, fetch or push is a git operation, which the git store
1692 /// charges g1t for: counted for billing once the answer has gone back
1693 /// (meters.rs), and a free workspace far past its share is slowed down
1694 /// rather than charged (see git_ops.rs). Whether it is past it is
1695 /// decided from counts this isolate already holds: the database is not
1696 /// asked on the way. A free workspace is never charged for private
1697 /// storage: once its private repositories hold the free amount, pushes
1698 /// to them stop, checked when a push begins so that git shows the
1699 /// reason. So do pushes to a repository at the store's size limit.
1700 async fn git_limits(
1701 &self,
1702 call: git_ops::GitCall,
1703 git: &git_http::GitRequest,
1704 repo: &Repo,
1705 env: &Env,
1706 ) -> Result<Option<(Response, u16, &'static str)>> {
1707 let namespace = git.path.namespace.to_lowercase();
1708 if meters::mapping_now().billable(call.meter()) > 0.0 {
1709 let now = now_ms();
1710 let hour = git_ops::hour_key(&rfc3339(now));
1711 let limits = git_ops::Limits::from_env(env);
1712 if let Some((month, hour_ops)) = git_ops::standing(&namespace, &hour, now)
1713 && git_ops::slow_down(month + 1, hour_ops + 1, limits.free_cap, limits.hourly)
1714 && git_ops::is_free_kept(env.service("BILLING").ok().as_ref(), &namespace).await
1715 {
1716 return Ok(Some((
1717 git_ops::too_many(&namespace, limits.free_cap, limits.hourly)?,
1718 429,
1719 "Too many git operations this hour.",
1720 )));
1721 }
1722 }
1723 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" {
1724 let held = self.held(repo).await;
1725 if held >= self.repo_limit {
1726 let message = format!(
1727 "{}/{} 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",
1728 repo.namespace,
1729 repo.name,
1730 pack_limits::megabytes(held)
1731 );
1732 return Ok(Some((Response::error(message, 403)?, 403, "The repository is full.")));
1733 }
1734 }
1735 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" && repo.is_private {
1736 let free = git_ops::free_private_bytes(env);
1737 let held = self.registry.private_bytes(&namespace).await.unwrap_or(0);
1738 if git_ops::storage_full(held, free)
1739 && git_ops::is_free(env.service("BILLING").ok().as_ref(), &namespace).await
1740 {
1741 return Ok(Some((
1742 git_ops::storage_full_response(&namespace, held, free)?,
1743 403,
1744 "Free private storage is full.",
1745 )));
1746 }
1747 }
1748 Ok(None)
1749 }
1750
1751 /// What a repository and its pull requests' working copies hold, as
1752 /// g1t counts it: read for a push's first request, kept a minute for
1753 /// the rest of it.
1754 async fn held(&self, repo: &Repo) -> u64 {
1755 let root = repo.fork_of.clone().unwrap_or_else(|| repo.id.clone());
1756 let now = now_ms();
1757 if let Some(held) = HELD.with(|held| held.borrow().get(&root, now)) {
1758 return held;
1759 }
1760 let held = self.registry.stored_bytes(&root).await.unwrap_or(0).max(0) as u64;
1761 HELD.with(|kept| kept.borrow_mut().put(root, held, now));
1762 held
1763 }
1764
1765 /// What a push changed, recorded once git has its answer.
1766 async fn record_push(&self, push: PushDone) -> Result<()> {
1767 let PushDone {
1768 repo,
1769 pushed,
1770 pack_bytes,
1771 actor,
1772 unscanned,
1773 } = push;
1774 // What the push stored, for billing's storage meter. A failure only
1775 // leaves the count short.
1776 if pack_bytes > 0
1777 && let Err(error) = self.registry.add_stored_bytes(&repo, pack_bytes).await
1778 {
1779 worker::console_error!("stored bytes for {} not counted: {error}", repo.name);
1780 }
1781 if pushed.is_empty() {
1782 return Ok(());
1783 }
1784 // Artifacts' own push notifications are per repository, which does
1785 // not fit a repo per pull request, so the front end reports pushes
1786 // itself: one event for each branch that moved.
1787 let stored = self.store.open(&store_key(&repo)).await?;
1788 for pushed in &pushed {
1789 // The store can refuse one ref and accept another, so each
1790 // branch is checked against where it actually is. A tag the
1791 // store cannot read back is taken as pushed.
1792 let moved = match pushed.branch() {
1793 Some(branch) => stored
1794 .log(branch, 1)
1795 .await?
1796 .first()
1797 .is_some_and(|commit| commit.hash == pushed.after),
1798 None => stored.log(&pushed.git_ref, 1).await.map_or(true, |head| {
1799 head.first().is_none_or(|commit| commit.hash == pushed.after)
1800 }),
1801 };
1802 if moved {
1803 self.publish_git_push(
1804 &repo,
1805 &pushed.git_ref,
1806 pushed.before.as_deref(),
1807 &pushed.after,
1808 actor.clone(),
1809 unscanned,
1810 )
1811 .await?;
1812 }
1813 }
1814 Ok(())
1815 }
1816}
1817
1818/// A push the store accepted, to be recorded once git has its answer.
1819struct PushDone {
1820 repo: Repo,
1821 pushed: Vec<git_http::Pushed>,
1822 pack_bytes: u64,
1823 actor: Option<String>,
1824 /// Too large to scan for secrets before it was stored.
1825 unscanned: bool,
1826}
1827
1828/// What a git request leaves for after its answer: its audit entry, with
1829/// how the request ended, and what a push changed.
1830struct AfterGit {
1831 audit: Option<Box<g1t_contracts::audit::NewAuditEntry>>,
1832 status: u16,
1833 message: Option<String>,
1834 push: Option<PushDone>,
1835}
1836
1837impl AfterGit {
1838 fn ended(&mut self, status: u16, message: Option<String>) {
1839 self.status = status;
1840 self.message = message;
1841 }
1842
1843 /// Does the work once the response is on its way. A failure is logged:
1844 /// git has already been told how its request went.
1845 fn spawn(self, env: &Env, ctx: &Context) {
1846 if self.audit.is_none() && self.push.is_none() {
1847 return;
1848 }
1849 let env = env.clone();
1850 ctx.wait_until(async move {
1851 let repos = match service(&env) {
1852 Ok(repos) => repos,
1853 Err(error) => {
1854 worker::console_error!("git request not recorded: {error}");
1855 return;
1856 }
1857 };
1858 repos.finish_git(self.audit, self.status, self.message).await;
1859 if let Some(push) = self.push
1860 && let Err(error) = repos.record_push(push).await
1861 {
1862 worker::console_error!("push not recorded: {error}");
1863 }
1864 });
1865 }
1866}
1867
1868fn service(env: &Env) -> Result<Repos<ArtifactsStore>> {
1869 let shared = shared::Shared::from_env(env).map(Rc::new);
1870 Ok(Repos {
1871 registry: Registry { db: env.d1("DB")? },
1872 store: ArtifactsStore::new(env, shared.clone())?,
1873 shared,
1874 events: env.service("EVENTS")?,
1875 security: env.service("SECURITY").ok(),
1876 billing: env.service("BILLING").ok(),
1877 identity: env.service("IDENTITY").ok(),
1878 free_private_bytes: git_ops::free_private_bytes(env),
1879 fork_days: forks::retention_days(env),
1880 repo_limit: env
1881 .var("REPO_STORAGE_LIMIT_BYTES")
1882 .ok()
1883 .and_then(|value| value.to_string().parse().ok())
1884 .unwrap_or(pack_limits::DEFAULT_REPO_LIMIT_BYTES),
1885 large_pushes: git_http::LargePushes::from_var(env.var("LARGE_PUSHES").ok().map(|value| value.to_string()).as_deref()),
1886 placement: shards::Placement::from_vars(
1887 env.var("ARTIFACTS_NEW_REPOS").ok().map(|value| value.to_string()).as_deref(),
1888 env.var("ARTIFACTS_EU_NAMESPACE").ok().map(|value| value.to_string()).as_deref(),
1889 ),
1890 })
1891}
1892
1893/// Writes what this isolate metered once the answer has gone back, every
1894/// few seconds at most (meters.rs).
1895fn flush_later(env: &Env, ctx: &Context) {
1896 if !meters::take_due() {
1897 return;
1898 }
1899 if let Ok(db) = env.d1("DB") {
1900 ctx.wait_until(async move { meters::flush(&db).await });
1901 }
1902}
1903
1904/// Read methods whose answer is an `Outcome`: when the git store is busy,
1905/// the site is told so in words instead of failing the page.
1906const OUTCOME_READS: [&str; 6] = ["tree", "blob", "log", "branches", "blame", "compare"];
1907
1908#[event(fetch)]
1909async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
1910 let mut repos = service(&env)?;
1911 let Some(method) = rpc_method(&request) else {
1912 let answered = repos.git_http(request, &env, &ctx).await;
1913 flush_later(&env, &ctx);
1914 return answered;
1915 };
1916 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
1917 // Git over HTTPS above always reads the primary.
1918 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
1919 repos.registry.db = db;
1920 let body: serde_json::Value = request.json().await?;
1921
1922 let answered = async { match method.as_str() {
1923 "get" => reply(&repos.get(args(body)?).await?),
1924 "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
1925 "readable" => {
1926 let a: ReadableArgs = args(body)?;
1927 reply(&repos.registry.readable(&a.ids, &a.viewer).await?)
1928 }
1929 "public_namespaces" => {
1930 let a: PublicNamespacesArgs = args(body)?;
1931 reply(&repos.registry.public_namespaces(&a.owner_id).await?)
1932 }
1933 "path_by_id" => {
1934 let a: PathByIdArgs = args(body)?;
1935 reply(
1936 &repos
1937 .registry
1938 .by_id(&a.id)
1939 .await?
1940 .filter(|repo| repo.fork_of.is_none())
1941 .map(|repo| RepoPath {
1942 namespace: repo.namespace,
1943 name: repo.name,
1944 }),
1945 )
1946 }
1947 "list" => {
1948 let a: ListArgs = args(body)?;
1949 reply(
1950 &repos
1951 .registry
1952 .list(
1953 &a.viewer,
1954 a.query.as_deref(),
1955 a.namespace.as_deref(),
1956 a.member_only,
1957 )
1958 .await?,
1959 )
1960 }
1961 "create" => reply(&repos.create(args(body)?).await?),
1962 // Services only: a GitHub mirror catching up, or pushing out.
1963 "mirror" => reply(&repos.mirror(args(body)?).await?),
1964 "transfer" => reply(&repos.transfer(args(body)?).await?),
1965 // A repository's lifecycle: see lifecycle.rs.
1966 "delete" => reply(&repos.delete(args(body)?).await?),
1967 "deleted" => reply(&repos.deleted(args(body)?).await?),
1968 "restore" => reply(&repos.restore(args(body)?).await?),
1969 "purge" => reply(&repos.purge(args(body)?).await?),
1970 "purge_due" => reply(&repos.purge_due(args(body)?).await?),
1971 "rename" => reply(&repos.rename(args(body)?).await?),
1972 "archive" => reply(&repos.archive(args(body)?).await?),
1973 "set_visibility" => reply(&repos.set_visibility(args(body)?).await?),
1974 "set_default_branch" => reply(&repos.set_default_branch(args(body)?).await?),
1975 "rename_branch" => reply(&repos.rename_branch(args(body)?).await?),
1976 "resolve_branch" => reply(&repos.resolve_branch(args(body)?).await?),
1977 "status_by_id" => reply(&repos.status_by_id(args(body)?).await?),
1978 "resolve_path" => {
1979 let a: ResolvePathArgs = args(body)?;
1980 reply(&repos.registry.resolve_moved(&a.path).await?)
1981 }
1982 "namespace_count" => {
1983 let a: NamespaceCountArgs = args(body)?;
1984 reply(&repos.registry.count_in(&a.namespace).await?)
1985 }
1986 "update" => reply(&repos.update(args(body)?).await?),
1987 "tree" => reply(&repos.tree(args(body)?).await?),
1988 "blob" => reply(&repos.blob(args(body)?).await?),
1989 "log" => reply(&repos.log(args(body)?).await?),
1990 "blame" => reply(&repos.blame(args(body)?).await?),
1991 "fork_for_pull" => reply(&repos.fork_for_pull(args(body)?).await?),
1992 "git_access" => reply(&repos.git_access(args(body)?).await?),
1993 "branches" => reply(&repos.branches(args(body)?).await?),
1994 "last_commits" => reply(&repos.last_commits(args(body)?).await?),
1995 "tags" => reply(&repos.tags(args(body)?).await?),
1996 "head" => reply(&repos.head(args(body)?).await?),
1997 "behind" => reply(&repos.behind(args(body)?).await?),
1998 "divergence" => reply(&repos.divergence(args(body)?).await?),
1999 "land" => reply(&repos.land(args(body)?).await?),
2000 "update_pull_branch" => reply(&repos.update_pull_branch(args(body)?).await?),
2001 "delete_branch" => reply(&repos.delete_branch(args(body)?).await?),
2002 "commit_file" => reply(&repos.commit_file(args(body)?).await?),
2003 "compare" => reply(&repos.compare(args(body)?).await?),
2004 "scan_history" => reply(&repos.scan_history(args(body)?).await?),
2005 "find_lockfiles" => reply(&repos.find_lockfiles(args(body)?).await?),
2006 "list_files" => reply(&repos.list_files(args(body)?).await?),
2007 "changed_files" => reply(&repos.changed_files(args(body)?).await?),
2008 "read_blobs" => reply(&repos.read_blobs(args(body)?).await?),
2009 // Services only: what the Composer registry builds packages from.
2010 "refs" => reply(&repos.refs_of(args(body)?).await?),
2011 "raw_file" => reply(&repos.raw_file(args(body)?).await?),
2012 "raw_blobs" => reply(&repos.raw_blobs(args(body)?).await?),
2013 "visibility" => {
2014 let a: g1t_contracts::repos::VisibilityArgs = args(body)?;
2015 reply(&repos.registry.visibility(&a.paths).await?)
2016 }
2017 "storage" => reply(&repos.registry.storage().await?),
2018 "git_operations" => {
2019 let a: GitOperationsArgs = args(body)?;
2020 reply(&git_ops::totals(&repos.registry.db, &a.month, a.since.as_deref(), a.namespace.as_deref().map(str::to_lowercase).as_deref()).await?)
2021 }
2022 "all_ids" => {
2023 let a: AllIdsArgs = args(body)?;
2024 let limit = a.limit.clamp(1, 500);
2025 let ids = repos.registry.ids_after(a.after.as_deref(), limit).await?;
2026 let next = (ids.len() == limit as usize).then(|| ids.last().cloned()).flatten();
2027 reply(&IdPage { ids, next })
2028 }
2029 // The raw meters of the git store, for reconciling with Cloudflare
2030 // (meters.rs, scripts/ops/artifacts-usage.mjs).
2031 "artifacts_usage" => {
2032 let a: meters::UsageArgs = args(body)?;
2033 reply(&meters::usage(&repos.registry.db, &a).await?)
2034 }
2035 "operation_mapping" => reply(&meters::read_mapping(&repos.registry.db).await?),
2036 // Services only: which meters are operations, changed without a deploy.
2037 "set_operation_mapping" => {
2038 let row: meters::MappingRow = args(body)?;
2039 meters::set_mapping(&repos.registry.db, &row, &rfc3339(now_ms())).await?;
2040 reply(&meters::read_mapping(&repos.registry.db).await?)
2041 }
2042 // How the git store has been answering, for the status page.
2043 "store_health" => {
2044 let a: meters::HealthArgs = args(body)?;
2045 reply(&meters::health(&repos.registry.db, &a).await?)
2046 }
2047 _ => Response::error("Unknown method", 404),
2048 } }
2049 .await;
2050 // The git store is busy: said in words, with when to try again.
2051 let answered = match answered {
2052 Err(error) => match resilience::busy(&error.to_string()) {
2053 Some(busy) if OUTCOME_READS.contains(&method.as_str()) => {
2054 reply(&Outcome::<()>::fail(FailureCode::Conflict, busy.message().trim()))
2055 }
2056 Some(busy) => {
2057 let response = Response::error(busy.message(), 503)?;
2058 response.headers().set("retry-after", &busy.retry_after.to_string())?;
2059 Ok(response)
2060 }
2061 None => Err(error),
2062 },
2063 answered => answered,
2064 };
2065 flush_later(&env, &ctx);
2066 served.finish(answered)
2067}
2068
2069/// The hourly sweep: deleted repositories whose time to be restored has
2070/// passed are purged. See lifecycle.rs.
2071#[event(scheduled)]
2072async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
2073 let repos = match service(&env) {
2074 Ok(repos) => repos,
2075 Err(error) => {
2076 worker::console_error!("repos: the sweep could not start: {error}");
2077 return;
2078 }
2079 };
2080 match repos.purge_due(PurgeDueArgs::default()).await {
2081 Ok(0) => {}
2082 Ok(count) => worker::console_log!("repos: purged {count} deleted repositories"),
2083 Err(error) => worker::console_error!("repos: the purge sweep failed: {error}"),
2084 }
2085 // Pull requests' working copies whose time has come (forks.rs).
2086 match repos.retire_due().await {
2087 Ok(0) => {}
2088 Ok(count) => worker::console_log!("repos: removed {count} pull request working copies"),
2089 Err(error) => worker::console_error!("repos: the working copy sweep failed: {error}"),
2090 }
2091 meters::flush(&repos.registry.db).await;
2092}
2093
2094/// Events from the bus. A workspace's rename: its repositories move to the
2095/// workspace's current slug, asked of identity by id, so a repeated or late
2096/// delivery lands in the same place; their git store keys stay as they
2097/// were. A workspace's deletion: its repositories are deleted with it,
2098/// restored with it, or purged with it.
2099#[event(queue)]
2100async fn queue(batch: MessageBatch<Event>, env: Env, ctx: Context) -> Result<()> {
2101 let registry = Registry { db: env.d1("DB")? };
2102 let identity = env.service("IDENTITY")?;
2103 let handled = handle_events(&batch, &env, &registry, &identity).await;
2104 flush_later(&env, &ctx);
2105 handled
2106}
2107
2108async fn handle_events(batch: &MessageBatch<Event>, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
2109 for message in batch.messages()? {
2110 let event = message.body();
2111 // A pull request merged, closed or reopened: its working copy is
2112 // kept or let go (forks.rs).
2113 if let Some(change) = forks::pull_change(&event.kind) {
2114 let Some(pull_id) = forks::pull_id_of(&event.data) else {
2115 worker::console_error!("{} {} names no pull request", event.kind, event.id);
2116 continue;
2117 };
2118 let repos = service(env)?;
2119 match change {
2120 forks::PullChange::Settled => repos.pull_settled(&pull_id).await?,
2121 forks::PullChange::Reopened => repos.pull_reopened(&pull_id).await?,
2122 }
2123 continue;
2124 }
2125 // A workspace deleted, restored or purged: its repositories go with
2126 // it, come back with it, or are purged with it (lifecycle.rs).
2127 if event.kind == "workspace.deleting" {
2128 match serde_json::from_value::<WorkspaceDeleting>(event.data.clone()) {
2129 Ok(deleting) => service(env)?.delete_with_workspace(&deleting, &protected_workspaces(env)).await?,
2130 Err(_) => worker::console_error!("workspace.deleting {} could not be read", event.id),
2131 }
2132 continue;
2133 }
2134 if event.kind == "workspace.restored" {
2135 match serde_json::from_value::<WorkspaceRestored>(event.data.clone()) {
2136 Ok(restored) => service(env)?.restore_with_workspace(&restored).await?,
2137 Err(_) => worker::console_error!("workspace.restored {} could not be read", event.id),
2138 }
2139 continue;
2140 }
2141 if event.kind == "workspace.deleted" {
2142 match serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) {
2143 Ok(deleted) => service(env)?.purge_workspace(&deleted, &protected_workspaces(env)).await?,
2144 Err(_) => worker::console_error!("workspace.deleted {} could not be read", event.id),
2145 }
2146 continue;
2147 }
2148 if event.kind != "workspace.renamed" {
2149 continue;
2150 }
2151 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
2152 worker::console_error!("workspace.renamed {} could not be read", event.id);
2153 continue;
2154 };
2155 let names: HashMap<String, String> = g1t_kit::call(
2156 identity,
2157 "usernames",
2158 &g1t_contracts::identity::UsernamesArgs {
2159 ids: vec![renamed.workspace_id.clone()],
2160 },
2161 )
2162 .await?;
2163 let current = names
2164 .get(&renamed.workspace_id)
2165 .cloned()
2166 .unwrap_or_else(|| renamed.to.clone());
2167 let left = registry
2168 .rename_namespace(&renamed.stale_slugs(&current), &current)
2169 .await?;
2170 if left > 0 {
2171 worker::console_error!(
2172 "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
2173 renamed.from,
2174 renamed.to
2175 );
2176 }
2177 }
2178 Ok(())
2179}
2180
2181/// The workspaces whose repositories never go with a deletion, whatever is
2182/// published: `PROTECTED_WORKSPACES` if set here, and Flagon's always.
2183fn protected_workspaces(env: &Env) -> Vec<String> {
2184 let configured = env.var("PROTECTED_WORKSPACES").ok().map(|v| v.to_string());
2185 g1t_contracts::identity::protected_names(configured.as_deref())
2186}
2187
2188/// The repository a push to a path that does not exist yet creates: private,
2189/// so nothing pushed by mistake is published. An owner makes it public on
2190/// purpose (`POST /repos/{owner}/{repo}/visibility`).
2191fn push_to_create(owner: &User, path: &RepoPath) -> CreateArgs {
2192 CreateArgs {
2193 owner: owner.clone(),
2194 namespace: path.namespace.clone(),
2195 name: path.name.clone(),
2196 description: None,
2197 is_private: true,
2198 import_url: None,
2199 import_token: None,
2200 }
2201}
2202
2203#[cfg(test)]
2204mod push_to_create_tests {
2205 use super::*;
2206
2207 #[test]
2208 fn a_pushed_repository_starts_private() {
2209 let owner: User = serde_json::from_value(serde_json::json!({ "id": "usr_1", "username": "ada" })).unwrap();
2210 let args = push_to_create(&owner, &RepoPath { namespace: "acme".into(), name: "site".into() });
2211 assert!(args.is_private);
2212 assert_eq!((args.namespace.as_str(), args.name.as_str()), ("acme", "site"));
2213 }
2214}