g1t/services/repos/src/lib.rs

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