g1t/services/repos/src/lib.rs

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