flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/repos/src/lib.rs

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