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,223 lines47,189 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 diff;
10mod git_http;
11mod import;
12mod land;
13mod listing;
14mod refs;
15mod registry;
16mod run_access;
17mod secret_scan;
18mod store;
19
20use g1t_contracts::events::{
21 Event, GitPush, NewEvent, Publish, RepoCreated, RepoForked, RepoUpdated, RepoVisibilityChanged,
22 WorkspaceRenamed,
23};
24use g1t_contracts::repos::*;
25use g1t_contracts::time::rfc3339;
26use g1t_contracts::{FailureCode, Outcome, PrincipalKind, User, Viewer, is_valid_repo_name, new_id};
27use g1t_kit::{args, now_ms, reply, rpc_method};
28use std::collections::{HashMap, HashSet, VecDeque};
29
30use serde::Serialize;
31use worker::{Context, Env, Fetcher, MessageBatch, Method, Request, Response, Result, event};
32
33use registry::{Registry, can_read, can_write, store_key};
34use store::{ArtifactsStore, GitRepo, GitStore, Scope};
35
36/// Namespace that holds every pull request's fork: `pulls/<pull id>`.
37const PULLS_NAMESPACE: &str = "pulls";
38const MAX_TEXT_BYTES: usize = 512 * 1024;
39/// How far back a pull request may have forked and still be landed.
40const MAX_ANCESTRY: u32 = 1000;
41const MAX_DESCRIPTION_CHARS: usize = 200;
42const SOURCE: &str = "repos";
43const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
44
45fn not_found<T>() -> Outcome<T> {
46 Outcome::fail(FailureCode::NotFound, "Repository not found.")
47}
48
49/// Decoded text, or `None` when the file is too large or looks binary.
50/// Whether a ref is a full commit hash rather than a branch name.
51fn is_commit_hash(git_ref: &str) -> bool {
52 git_ref.len() == 40 && git_ref.bytes().all(|b| b.is_ascii_hexdigit())
53}
54
55fn text_of(bytes: Vec<u8>) -> Option<String> {
56 if bytes.len() > MAX_TEXT_BYTES || bytes.contains(&0) {
57 return None;
58 }
59 Some(String::from_utf8_lossy(&bytes).into_owned())
60}
61
62fn is_readme(name: &str) -> bool {
63 matches!(
64 name.to_lowercase().as_str(),
65 "readme" | "readme.md" | "readme.markdown" | "readme.txt"
66 )
67}
68
69/// Whether `ancestor` is reachable from the newest commit in `history`.
70///
71/// `history` is the first-parent chain, which is all the store lists; a fork
72/// that merged the target branch in has the target's head on a second
73/// parent, so the walk follows every parent.
74async fn descends_from<R: GitRepo>(repo: &R, history: &[Commit], ancestor: &str) -> Result<bool> {
75 let known: HashMap<&str, &[String]> = history
76 .iter()
77 .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
78 .collect();
79 let mut seen = HashSet::new();
80 let mut queue: Vec<String> = history
81 .first()
82 .map(|c| c.hash.clone())
83 .into_iter()
84 .collect();
85 while let Some(hash) = queue.pop() {
86 if hash == ancestor {
87 return Ok(true);
88 }
89 if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
90 continue;
91 }
92 match known.get(hash.as_str()) {
93 Some(parents) => queue.extend(parents.iter().cloned()),
94 None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
95 }
96 }
97 Ok(false)
98}
99
100/// The commit closest to the newest in `history` that is also in `shared`:
101/// where a fork and the repository it came from last agreed.
102async fn nearest_ancestor_in<R: GitRepo>(
103 repo: &R,
104 history: &[Commit],
105 shared: &HashSet<String>,
106) -> Result<Option<String>> {
107 let known: HashMap<&str, &[String]> = history
108 .iter()
109 .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
110 .collect();
111 let mut seen = HashSet::new();
112 let mut queue: VecDeque<String> = history
113 .first()
114 .map(|c| c.hash.clone())
115 .into_iter()
116 .collect();
117 while let Some(hash) = queue.pop_front() {
118 if shared.contains(&hash) {
119 return Ok(Some(hash));
120 }
121 if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
122 continue;
123 }
124 match known.get(hash.as_str()) {
125 Some(parents) => queue.extend(parents.iter().cloned()),
126 None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
127 }
128 }
129 Ok(None)
130}
131
132struct Repos<S: GitStore> {
133 registry: Registry,
134 store: S,
135 events: Fetcher,
136 /// Asked during a push which secrets have been allowed.
137 security: Option<Fetcher>,
138}
139
140impl<S: GitStore> Repos<S> {
141 async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> {
142 g1t_kit::call(
143 &self.events,
144 "publish",
145 &Publish {
146 events: vec![event],
147 },
148 )
149 .await
150 }
151
152 /// Whether the viewer may read `repo`. A pull request's fork of a
153 /// private repository can be read by everyone who can read that
154 /// repository, so its members can review and check out the change, as
155 /// well as by whoever opened the pull request.
156 async fn may_read(&self, repo: &Repo, viewer: &Viewer) -> Result<bool> {
157 if can_read(repo, viewer) {
158 return Ok(true);
159 }
160 let Some(source_id) = &repo.fork_of else {
161 return Ok(false);
162 };
163 Ok(self
164 .registry
165 .by_id(source_id)
166 .await?
167 .is_some_and(|source| can_read(&source, viewer)))
168 }
169
170 /// `repo`, if there is one and the viewer may read it.
171 async fn visible(&self, repo: Option<Repo>, viewer: &Viewer) -> Result<Option<Repo>> {
172 Ok(match repo {
173 Some(repo) if self.may_read(&repo, viewer).await? => Some(repo),
174 _ => None,
175 })
176 }
177
178 /// Resolves a repo the viewer may read; private repos look missing.
179 async fn readable(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
180 self.visible(self.registry.by_path(path).await?, viewer)
181 .await
182 }
183
184 async fn get(&self, a: GetArgs) -> Result<Outcome<Repo>> {
185 Ok(self
186 .readable(&a.path, &a.viewer)
187 .await?
188 .map_or_else(not_found, Outcome::Ok))
189 }
190
191 async fn get_by_id(&self, a: GetByIdArgs) -> Result<Outcome<Repo>> {
192 Ok(self
193 .visible(self.registry.by_id(&a.id).await?, &a.viewer)
194 .await?
195 .map_or_else(not_found, Outcome::Ok))
196 }
197
198 async fn update(&self, a: UpdateArgs) -> Result<Outcome<Repo>> {
199 let viewer = Some(a.actor.clone());
200 let Some(repo) = self.readable(&a.path, &viewer).await? else {
201 return Ok(not_found());
202 };
203 if repo.fork_of.is_some() || !can_write(&repo, &viewer) {
204 return Ok(Outcome::fail(
205 FailureCode::Forbidden,
206 "Only members of the repository's workspace can change its settings.",
207 ));
208 }
209 if !a.actor.verified {
210 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
211 }
212 let description = match a.description {
213 Some(text) => Some(
214 text.trim()
215 .chars()
216 .take(MAX_DESCRIPTION_CHARS)
217 .collect::<String>(),
218 )
219 .filter(|text| !text.is_empty()),
220 None => repo.description.clone(),
221 };
222 let is_private = a.is_private.unwrap_or(repo.is_private);
223 let protected = a.protected.unwrap_or(repo.protected);
224 let topics = match &a.topics {
225 Some(topics) => match clean_topics(topics) {
226 Ok(topics) => topics,
227 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
228 },
229 None => repo.topics.clone(),
230 };
231 self.registry
232 .update(&repo.id, description.as_deref(), is_private, protected, &topics)
233 .await?;
234 let visibility_changed = is_private != repo.is_private;
235 let updated = Repo {
236 description,
237 is_private,
238 protected,
239 topics,
240 ..repo
241 };
242 // Search and anything else that shows the repository hears of it;
243 // a change of visibility is announced on its own as well, so that
244 // what was public stops being shown at once.
245 self.publish(NewEvent {
246 kind: "repo.updated",
247 source: SOURCE,
248 repo_id: Some(updated.id.clone()),
249 actor: Some(a.actor.id.clone()),
250 data: RepoUpdated {
251 repo_id: updated.id.clone(),
252 namespace: updated.namespace.clone(),
253 name: updated.name.clone(),
254 is_private,
255 visibility_changed,
256 },
257 })
258 .await?;
259 if visibility_changed {
260 self.publish(NewEvent {
261 kind: "repo.visibility_changed",
262 source: SOURCE,
263 repo_id: Some(updated.id.clone()),
264 actor: Some(a.actor.id),
265 data: RepoVisibilityChanged {
266 repo_id: updated.id.clone(),
267 is_private,
268 },
269 })
270 .await?;
271 }
272 Ok(Outcome::Ok(updated))
273 }
274
275 /// The repository with this id, if it is not a fork, and its store.
276 async fn stored(&self, repo_id: &str) -> Result<Option<S::Repo>> {
277 match self.registry.by_id(repo_id).await? {
278 Some(repo) if repo.fork_of.is_none() => Ok(Some(self.store.open(&store_key(&repo)).await?)),
279 _ => Ok(None),
280 }
281 }
282
283 async fn list_files(&self, a: ListFilesArgs) -> Result<FileList> {
284 let Some(repo) = self.registry.by_id(&a.repo_id).await?.filter(|repo| repo.fork_of.is_none()) else {
285 return Ok(FileList::default());
286 };
287 let git = self.store.open(&store_key(&repo)).await?;
288 let head = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
289 listing::list(&git, None, &head, &a.skip_dirs, a.limit).await
290 }
291
292 async fn changed_files(&self, a: ChangedFilesArgs) -> Result<FileList> {
293 let Some(git) = self.stored(&a.repo_id).await? else {
294 return Ok(FileList::default());
295 };
296 listing::list(&git, a.base.as_deref(), &a.head, &a.skip_dirs, a.limit).await
297 }
298
299 async fn read_blobs(&self, a: ReadBlobsArgs) -> Result<Vec<BlobText>> {
300 let Some(git) = self.stored(&a.repo_id).await? else {
301 return Ok(Vec::new());
302 };
303 listing::read(&git, &a.hashes, a.max_bytes.min(MAX_TEXT_BYTES as u32)).await
304 }
305
306 async fn create(&self, a: CreateArgs) -> Result<Outcome<Repo>> {
307 if !a.owner.verified {
308 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
309 }
310 let name = a.name.trim().to_lowercase();
311 if !is_valid_repo_name(&name) {
312 return Ok(Outcome::fail(
313 FailureCode::Invalid,
314 "Use letters, digits, dots, hyphens and underscores only.",
315 ));
316 }
317 let namespace = a.namespace.trim().to_lowercase();
318 if namespace.is_empty() {
319 return Ok(Outcome::fail(
320 FailureCode::Invalid,
321 "Say which workspace to create the repository in.",
322 ));
323 }
324 if !a.owner.is_member(&namespace) {
325 return Ok(Outcome::fail(
326 FailureCode::Forbidden,
327 "You are not a member of that workspace.",
328 ));
329 }
330 let path = RepoPath { namespace, name };
331 if self.registry.by_path(&path).await?.is_some() {
332 return Ok(Outcome::fail(
333 FailureCode::Conflict,
334 "That workspace already has a repository with that name.",
335 ));
336 }
337 // An import is fetched before anything is created, so that an
338 // address that does not work leaves nothing behind.
339 let mut imported = None;
340 if let Some(url) = a
341 .import_url
342 .as_deref()
343 .map(str::trim)
344 .filter(|url| !url.is_empty())
345 {
346 let Some(url) = import::clean_url(url) else {
347 return Ok(Outcome::fail(
348 FailureCode::Invalid,
349 "Give the https address of a public repository, such as https://github.com/owner/repo.",
350 ));
351 };
352 let remote = match import::discover(&url).await? {
353 Ok(remote) => remote,
354 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
355 };
356 let pack = match import::fetch(&url, &remote.head).await? {
357 Ok(pack) => pack,
358 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
359 };
360 imported = Some((remote, pack));
361 }
362 let now = now_ms();
363 let repo = Repo {
364 id: new_id("rep", now),
365 namespace: path.namespace,
366 name: path.name,
367 description: a
368 .description
369 .map(|text| text.trim().to_owned())
370 .filter(|text| !text.is_empty()),
371 is_private: a.is_private,
372 owner_id: a.owner.id.clone(),
373 default_branch: imported
374 .as_ref()
375 .map_or_else(|| "main".to_owned(), |(remote, _)| remote.branch.clone()),
376 fork_of: None,
377 protected: false,
378 created_at: rfc3339(now),
379 topics: Vec::new(),
380 };
381 self.registry.claim_store_key(&repo).await?;
382 self.store
383 .create(
384 &store_key(&repo),
385 repo.description.as_deref(),
386 &repo.default_branch,
387 )
388 .await?;
389 self.registry.insert(&repo).await?;
390 let mut pushed = None;
391 if let Some((remote, pack)) = imported {
392 let access = self
393 .store
394 .open(&store_key(&repo))
395 .await?
396 .access(Scope::Write)
397 .await?;
398 let stored =
399 land::push_pack(&access, &repo.default_branch, None, &remote.head, pack).await?;
400 if let Err(reason) = stored {
401 self.registry.remove(&repo.id).await?;
402 return Ok(Outcome::fail(
403 FailureCode::Invalid,
404 format!("The repository could not be stored: {reason}"),
405 ));
406 }
407 pushed = Some(remote.head);
408 }
409 self.publish(NewEvent {
410 kind: "repo.created",
411 source: SOURCE,
412 repo_id: Some(repo.id.clone()),
413 actor: Some(a.owner.id),
414 data: RepoCreated {
415 repo_id: repo.id.clone(),
416 namespace: repo.namespace.clone(),
417 name: repo.name.clone(),
418 is_private: repo.is_private,
419 },
420 })
421 .await?;
422 if let Some(head) = pushed {
423 self.publish_push(
424 &repo,
425 &format!("refs/heads/{}", repo.default_branch),
426 None,
427 &head,
428 None,
429 )
430 .await?;
431 }
432 Ok(Outcome::Ok(repo))
433 }
434
435 async fn tree(&self, a: TreeArgs) -> Result<Outcome<TreeView>> {
436 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
437 return Ok(not_found());
438 };
439 let git = self.store.open(&store_key(&repo)).await?;
440 let git_ref = a
441 .git_ref
442 .clone()
443 .unwrap_or_else(|| repo.default_branch.clone());
444
445 let Some(head) = git.log(&git_ref, 1).await?.into_iter().next() else {
446 // An unknown ref is an error; a repo with no commits is just empty.
447 if a.git_ref.is_some() {
448 return Ok(Outcome::fail(
449 FailureCode::NotFound,
450 "No such branch, tag or commit.",
451 ));
452 }
453 return Ok(Outcome::Ok(TreeView {
454 repo,
455 git_ref,
456 path: a.tree_path,
457 head: None,
458 entries: Vec::new(),
459 readme: None,
460 }));
461 };
462
463 let no_directory = || Outcome::fail(FailureCode::NotFound, "No such directory.");
464 let mut entries = git.read_tree(&head.tree_hash).await?;
465 for segment in a.tree_path.split('/').filter(|segment| !segment.is_empty()) {
466 let next = entries.as_ref().and_then(|entries| {
467 entries
468 .iter()
469 .find(|entry| entry.name == segment && entry.kind == EntryKind::Tree)
470 });
471 let Some(next) = next else {
472 return Ok(no_directory());
473 };
474 entries = git.read_tree(&next.hash).await?;
475 }
476 let Some(mut entries) = entries else {
477 return Ok(no_directory());
478 };
479 // Directories first, then by name.
480 entries.sort_by(|a, b| {
481 (b.kind == EntryKind::Tree)
482 .cmp(&(a.kind == EntryKind::Tree))
483 .then_with(|| a.name.cmp(&b.name))
484 });
485
486 let readme_entry = entries
487 .iter()
488 .find(|entry| entry.kind == EntryKind::Blob && is_readme(&entry.name));
489 let readme = match readme_entry {
490 Some(entry) => git.read_blob(&entry.hash).await?.map(|bytes| Readme {
491 name: entry.name.clone(),
492 text: text_of(bytes),
493 }),
494 None => None,
495 };
496 Ok(Outcome::Ok(TreeView {
497 repo,
498 git_ref,
499 path: a.tree_path,
500 head: Some(head),
501 entries,
502 readme,
503 }))
504 }
505
506 async fn blob(&self, a: BlobArgs) -> Result<Outcome<BlobView>> {
507 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
508 return Ok(not_found());
509 };
510 let bytes = if a.file_path.is_empty() {
511 None
512 } else {
513 let git = self.store.open(&store_key(&repo)).await?;
514 git.read_file(&a.git_ref, &a.file_path).await?
515 };
516 let Some(bytes) = bytes else {
517 return Ok(Outcome::fail(FailureCode::NotFound, "No such file."));
518 };
519 Ok(Outcome::Ok(BlobView {
520 repo,
521 git_ref: a.git_ref,
522 path: a.file_path,
523 size: bytes.len() as u64,
524 text: text_of(bytes),
525 }))
526 }
527
528 async fn blame(&self, a: BlameArgs) -> Result<Outcome<Blame>> {
529 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
530 return Ok(not_found());
531 };
532 let git = self.store.open(&store_key(&repo)).await?;
533 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
534 Ok(match blame::blame(&git, &git_ref, &a.file_path).await? {
535 Some(blame) => Outcome::Ok(blame),
536 None => not_found(),
537 })
538 }
539
540 async fn log(&self, a: LogArgs) -> Result<Outcome<Vec<Commit>>> {
541 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
542 return Ok(not_found());
543 };
544 let git = self.store.open(&store_key(&repo)).await?;
545 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
546 Ok(Outcome::Ok(git.log(&git_ref, a.limit).await?))
547 }
548
549 /// The repository's branches, default branch first.
550 async fn branches(&self, a: BranchesArgs) -> Result<Outcome<Vec<Branch>>> {
551 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
552 return Ok(not_found());
553 };
554 let mut branches = self.store.open(&store_key(&repo)).await?.branches().await?;
555 branches.sort_by_key(|branch| branch.name != repo.default_branch);
556 Ok(Outcome::Ok(branches))
557 }
558
559 /// Whether a pull request's source lacks commits that the branch it
560 /// would merge into has.
561 async fn behind(&self, a: BehindArgs) -> Result<bool> {
562 let Some(source) = self.registry.by_id(&a.source_id).await? else {
563 return Ok(false);
564 };
565 let target = match &source.fork_of {
566 Some(id) => self.registry.by_id(id).await?,
567 None => Some(source.clone()),
568 };
569 let Some(target) = target else {
570 return Ok(false);
571 };
572 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
573 let target_head = self
574 .store
575 .open(&store_key(&target))
576 .await?
577 .log(&target.default_branch, 1)
578 .await?
579 .into_iter()
580 .next()
581 .map(|commit| commit.hash);
582 let Some(target_head) = target_head else {
583 return Ok(false);
584 };
585 let source_git = self.store.open(&store_key(&source)).await?;
586 let history = source_git.log(&branch, MAX_ANCESTRY).await?;
587 if history.is_empty() {
588 return Ok(false);
589 }
590 Ok(!descends_from(&source_git, &history, &target_head).await?)
591 }
592
593 /// The files a pull request's source and the default branch it would
594 /// merge into each changed since they last agreed. Where the two lists
595 /// share no file, the merge cannot conflict; where they do, it may.
596 async fn divergence(&self, a: BehindArgs) -> Result<Option<Divergence>> {
597 let Some(source) = self.registry.by_id(&a.source_id).await? else {
598 return Ok(None);
599 };
600 let target = match &source.fork_of {
601 Some(id) => self.registry.by_id(id).await?,
602 None => Some(source.clone()),
603 };
604 let Some(target) = target else {
605 return Ok(None);
606 };
607 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
608 let source_git = self.store.open(&store_key(&source)).await?;
609 let target_git = self.store.open(&store_key(&target)).await?;
610 let (history, target_history) = futures_util::future::try_join(
611 source_git.log(&branch, MAX_ANCESTRY),
612 target_git.log(&target.default_branch, MAX_ANCESTRY),
613 )
614 .await?;
615 let (Some(head), Some(base)) = (history.first(), target_history.first()) else {
616 return Ok(None);
617 };
618 let behind = !descends_from(&source_git, &history, &base.hash).await?;
619 let shared: HashSet<String> = target_history.iter().map(|commit| commit.hash.clone()).collect();
620 let merge_base = nearest_ancestor_in(&source_git, &history, &shared).await?;
621 let mut divergence = Divergence {
622 head: head.hash.clone(),
623 base: base.hash.clone(),
624 merge_base: merge_base.clone(),
625 behind,
626 ..Divergence::default()
627 };
628 let merge_base_tree = match &merge_base {
629 Some(hash) => target_history
630 .iter()
631 .find(|commit| commit.hash == *hash)
632 .map(|commit| commit.tree_hash.clone()),
633 None => None,
634 };
635 let Some(merge_base_tree) = merge_base_tree else {
636 // No common history to compare from: say nothing is known.
637 divergence.truncated = true;
638 return Ok(Some(divergence));
639 };
640 let (ours, truncated_ours) =
641 diff::changed_paths(&source_git, Some(&merge_base_tree), &head.tree_hash).await?;
642 divergence.ours = ours;
643 divergence.truncated = truncated_ours;
644 if behind {
645 let (theirs, truncated_theirs) =
646 diff::changed_paths(&target_git, Some(&merge_base_tree), &base.tree_hash).await?;
647 divergence.theirs = theirs;
648 divergence.truncated |= truncated_theirs;
649 }
650 Ok(Some(divergence))
651 }
652
653 async fn head(&self, a: HeadArgs) -> Result<Option<String>> {
654 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
655 return Ok(None);
656 };
657 let branch = if a.branch.is_empty() { &repo.default_branch } else { &a.branch };
658 let git = self.store.open(&store_key(&repo)).await?;
659 Ok(git
660 .log(branch, 1)
661 .await?
662 .into_iter()
663 .next()
664 .map(|commit| commit.hash))
665 }
666
667 async fn delete_branch(&self, a: DeleteBranchArgs) -> Result<Outcome<bool>> {
668 if !a.branch.starts_with(G1T_BRANCH_PREFIX) {
669 return Ok(Outcome::fail(
670 FailureCode::Forbidden,
671 "Only branches g1t made for itself can be deleted this way.",
672 ));
673 }
674 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
675 return Ok(not_found());
676 };
677 let git = self.store.open(&store_key(&repo)).await?;
678 let Some(old) = git
679 .branches()
680 .await?
681 .into_iter()
682 .find(|branch| branch.name == a.branch)
683 .map(|branch| branch.hash)
684 else {
685 return Ok(Outcome::Ok(false));
686 };
687 let access = git.access(Scope::Write).await?;
688 if let Err(reason) = land::delete_ref(&access, &a.branch, &old).await? {
689 return Ok(Outcome::fail(
690 FailureCode::Conflict,
691 format!("{} could not be deleted: {reason}", a.branch),
692 ));
693 }
694 Ok(Outcome::Ok(true))
695 }
696
697 async fn fork_for_pull(&self, a: ForkArgs) -> Result<Outcome<Repo>> {
698 let viewer = Some(a.actor.clone());
699 let Some(source) = self
700 .registry
701 .by_id(&a.source_id)
702 .await?
703 .filter(|repo| can_read(repo, &viewer))
704 else {
705 return Ok(not_found());
706 };
707 let now = now_ms();
708 let fork = Repo {
709 id: new_id("rep", now),
710 namespace: PULLS_NAMESPACE.to_owned(),
711 name: a.pull_id.clone(),
712 description: None,
713 // A fork is exactly as visible as the repo it came from.
714 is_private: source.is_private,
715 owner_id: a.actor.id.clone(),
716 default_branch: source.default_branch.clone(),
717 fork_of: Some(source.id.clone()),
718 protected: false,
719 created_at: rfc3339(now),
720 topics: Vec::new(),
721 };
722 self.registry.claim_store_key(&fork).await?;
723 self.store
724 .open(&store_key(&source))
725 .await?
726 .fork(&store_key(&fork))
727 .await?;
728 self.registry.insert(&fork).await?;
729 self.publish(NewEvent {
730 kind: "repo.forked",
731 source: SOURCE,
732 repo_id: Some(source.id.clone()),
733 actor: Some(a.actor.id),
734 data: RepoForked {
735 repo_id: fork.id.clone(),
736 source_repo_id: source.id,
737 pull_id: a.pull_id,
738 },
739 })
740 .await?;
741 Ok(Outcome::Ok(fork))
742 }
743
744 async fn git_access(&self, a: GitAccessArgs) -> Result<Outcome<GitAccess>> {
745 let write = a.service == GitService::ReceivePack;
746 // Anonymous callers are asked to authenticate whether or not the repo
747 // exists, so private repos cannot be told apart from missing ones.
748 let denied = || match &a.viewer {
749 Some(_) => not_found(),
750 None => Outcome::fail(FailureCode::Unauthenticated, "Authentication required."),
751 };
752 // An agent's token works through the API only: its sandbox has its
753 // own way to push, to its own pull request.
754 if a.viewer.as_ref().is_some_and(|user| user.kind == PrincipalKind::Agent) {
755 return Ok(Outcome::fail(
756 FailureCode::Forbidden,
757 "A g1t agent's token cannot be used with git.",
758 ));
759 }
760 if let (true, Some(user)) = (write, &a.viewer)
761 && !user.verified
762 {
763 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
764 }
765
766 let repo = match self.registry.by_path(&a.path).await? {
767 Some(repo) => {
768 let allowed = if write {
769 can_write(&repo, &a.viewer)
770 } else {
771 self.may_read(&repo, &a.viewer).await?
772 };
773 if !allowed {
774 return Ok(denied());
775 }
776 repo
777 }
778 None => {
779 // Push to create, in a workspace the pusher belongs to.
780 let owner = a
781 .viewer
782 .as_ref()
783 .filter(|user| write && user.is_member(&a.path.namespace.to_lowercase()));
784 let Some(owner) = owner else {
785 return Ok(denied());
786 };
787 let created = self
788 .create(CreateArgs {
789 owner: owner.clone(),
790 namespace: a.path.namespace.clone(),
791 name: a.path.name.clone(),
792 description: None,
793 is_private: false,
794 import_url: None,
795 })
796 .await?;
797 match created {
798 Outcome::Ok(repo) => repo,
799 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
800 }
801 }
802 };
803 let git = self.store.open(&store_key(&repo)).await?;
804 let scope = if write { Scope::Write } else { Scope::Read };
805 Ok(Outcome::Ok(git.access(scope).await?))
806 }
807
808 async fn land(&self, a: LandArgs) -> Result<Outcome<Landed>> {
809 let actor: Viewer = Some(a.actor.clone());
810 let Some(source) = self.registry.by_id(&a.source_id).await? else {
811 return Ok(not_found());
812 };
813 // A fork lands on the repository it came from; a branch on its own.
814 let target = match &source.fork_of {
815 Some(id) => self.registry.by_id(id).await?,
816 None => Some(source.clone()),
817 };
818 let Some(target) = target.filter(|repo| can_read(repo, &actor)) else {
819 return Ok(not_found());
820 };
821 if !can_write(&target, &actor) {
822 return Ok(Outcome::fail(
823 FailureCode::Forbidden,
824 "Only members of the repository's workspace can merge a pull request.",
825 ));
826 }
827 if !a.actor.verified {
828 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
829 }
830
831 let branch = &target.default_branch;
832 let from_fork = source.id != target.id;
833 let source_branch = match a.branch {
834 Some(name) if !from_fork && name == *branch => {
835 return Ok(Outcome::fail(
836 FailureCode::Invalid,
837 format!("{branch} cannot be merged into itself."),
838 ));
839 }
840 Some(name) => name,
841 None if from_fork => branch.clone(),
842 None => {
843 return Ok(Outcome::fail(
844 FailureCode::Invalid,
845 "Say which branch to merge.",
846 ));
847 }
848 };
849
850 let source_git = self.store.open(&store_key(&source)).await?;
851 let target_git = self.store.open(&store_key(&target)).await?;
852 let history = source_git.log(&source_branch, MAX_ANCESTRY).await?;
853 let Some(new) = history.first().map(|commit| commit.hash.clone()) else {
854 return Ok(Outcome::fail(
855 FailureCode::Conflict,
856 "This pull request has no commits to merge.",
857 ));
858 };
859 let old = target_git
860 .log(branch, 1)
861 .await?
862 .into_iter()
863 .next()
864 .map(|commit| commit.hash);
865
866 if old.as_deref() == Some(new.as_str()) {
867 return Ok(Outcome::Ok(Landed {
868 commit: new,
869 previous: None,
870 }));
871 }
872 // Moving the branch to a commit that does not descend from its
873 // current head would discard whatever landed in between.
874 if let Some(old) = &old
875 && !descends_from(&source_git, &history, old).await?
876 {
877 let remedy = if from_fork {
878 format!("Pull {branch} into the pull request's fork, push, and merge again.")
879 } else {
880 format!("Merge {branch} into {source_branch}, push, and merge again.")
881 };
882 return Ok(Outcome::fail(
883 FailureCode::Conflict,
884 format!("{branch} has moved since this pull request was opened. {remedy}"),
885 ));
886 }
887
888 // For a branch the objects are already in the target; sending them
889 // again is harmless and keeps one way of moving a ref.
890 let source_access = source_git.access(Scope::Read).await?;
891 let target_access = target_git.access(Scope::Write).await?;
892 let pushed =
893 land::fast_forward(&source_access, &target_access, branch, old.as_deref(), &new)
894 .await?;
895 if let Err(reason) = pushed {
896 // Most often another pull request landed between the check and the push.
897 return Ok(Outcome::fail(
898 FailureCode::Conflict,
899 format!("{branch} could not be updated: {reason}"),
900 ));
901 }
902 self.publish_push(
903 &target,
904 &format!("refs/heads/{branch}"),
905 old.as_deref(),
906 &new,
907 Some(a.actor.id),
908 )
909 .await?;
910 Ok(Outcome::Ok(Landed {
911 commit: new,
912 previous: old,
913 }))
914 }
915
916 async fn compare(&self, a: CompareArgs) -> Result<Outcome<Comparison>> {
917 let Some(repo) = self
918 .visible(self.registry.by_id(&a.repo_id).await?, &a.viewer)
919 .await?
920 else {
921 return Ok(not_found());
922 };
923 let git = self.store.open(&store_key(&repo)).await?;
924 let head_ref = a.head.as_deref().unwrap_or(&repo.default_branch);
925 // The head's history is only searched when the base is worked out
926 // from another branch.
927 let depth = if a.base.is_some() || is_commit_hash(head_ref) { 1 } else { MAX_ANCESTRY };
928 let history = git.log(head_ref, depth).await?;
929 let Some(head) = history.first() else {
930 return Ok(Outcome::fail(
931 FailureCode::Conflict,
932 "There are no commits to compare.",
933 ));
934 };
935
936 // Where the head's history meets the default branch of `against`.
937 let shared_with = async |against: &Repo| -> Result<Option<String>> {
938 let against_git = self.store.open(&store_key(against)).await?;
939 let shared: HashSet<String> = against_git
940 .log(&against.default_branch, MAX_ANCESTRY)
941 .await?
942 .into_iter()
943 .map(|commit| commit.hash)
944 .collect();
945 nearest_ancestor_in(&git, &history, &shared).await
946 };
947 let base = match (a.base, &repo.fork_of) {
948 (Some(base), _) => Some(base),
949 // A fork is compared with the last commit it shares with the
950 // repository it came from.
951 (None, Some(target_id)) => match self.registry.by_id(target_id).await? {
952 Some(target) => shared_with(&target).await?,
953 None => None,
954 },
955 // A branch, with the point where it left the default branch.
956 // A single commit, with its first parent.
957 (None, None) if is_commit_hash(head_ref) => head.parents.first().cloned(),
958 (None, None) if head_ref != repo.default_branch => shared_with(&repo).await?,
959 (None, None) => head.parents.first().cloned(),
960 };
961 let base_tree = match &base {
962 Some(base) => git
963 .log(base, 1)
964 .await?
965 .into_iter()
966 .next()
967 .map(|commit| commit.tree_hash),
968 None => None,
969 };
970 let (files, truncated) =
971 diff::compare_trees(&git, base_tree.as_deref(), &head.tree_hash).await?;
972 Ok(Outcome::Ok(Comparison {
973 base,
974 head: head.hash.clone(),
975 files,
976 truncated,
977 }))
978 }
979
980 /// Reports that `git_ref` of `repo` (a full ref) now points to `after`.
981 async fn publish_push(
982 &self,
983 repo: &Repo,
984 git_ref: &str,
985 before: Option<&str>,
986 after: &str,
987 actor: Option<String>,
988 ) -> Result<()> {
989 self.publish(NewEvent {
990 kind: "git.push",
991 source: SOURCE,
992 repo_id: Some(repo.id.clone()),
993 actor,
994 data: GitPush {
995 repo_id: repo.id.clone(),
996 git_ref: git_ref.to_owned(),
997 before: before.map(str::to_owned),
998 after: after.to_owned(),
999 default_branch: git_ref.strip_prefix("refs/heads/")
1000 == Some(repo.default_branch.as_str()),
1001 },
1002 })
1003 .await
1004 }
1005
1006 /// Git over HTTPS.
1007 async fn git_http(&self, request: Request, env: &Env) -> Result<Response> {
1008 let Some(git) = git_http::parse(&request.url()?) else {
1009 return Response::error("Not found", 404);
1010 };
1011 // A workspace that was renamed: git follows a redirect when it
1012 // first asks for refs, and uses the new address from then on.
1013 if self.registry.by_path(&git.path).await?.is_none()
1014 && let Some(location) = git_http::renamed(&request.url()?, &env.service("IDENTITY")?).await?
1015 {
1016 return git_http::moved(&location, request.method() == Method::Get);
1017 }
1018 let viewer = git_http::viewer(&request, &env.service("IDENTITY")?).await?;
1019 // A run credential is checked against its grants, then acts as the
1020 // person it works for. See run_access.rs.
1021 let (request, viewer, audit) = match self.admit_git(request, &git, viewer).await? {
1022 run_access::Admitted::Go { request, viewer, entry } => (request, viewer, entry),
1023 run_access::Admitted::Refused(response) => return Ok(response),
1024 };
1025 let access = self
1026 .git_access(GitAccessArgs {
1027 path: git.path.clone(),
1028 viewer: viewer.clone(),
1029 service: git.service,
1030 })
1031 .await?;
1032 let access = match access {
1033 Outcome::Ok(access) => access,
1034 refused => {
1035 let response = git_http::refuse(refused)?;
1036 self.finish_git(audit, response.status_code(), None).await;
1037 return Ok(response);
1038 }
1039 };
1040 // A protected default branch takes changes only from a merged pull
1041 // request, which lands without going through here.
1042 let protected = match self.registry.by_path(&git.path).await? {
1043 Some(repo) if repo.protected && repo.fork_of.is_none() => Some(repo.default_branch),
1044 _ => None,
1045 };
1046 // Push protection: a push that adds a secret is refused. See secret_scan.rs.
1047 let scan = async |body: &[u8]| self.protect(&git.path, viewer.as_ref(), body).await;
1048 let forwarded =
1049 match git_http::forward(request, &git, &access, protected.as_deref(), scan).await? {
1050 git_http::Push::Forwarded(forwarded) => forwarded,
1051 git_http::Push::Refused(response) => {
1052 self.finish_git(audit, 403, Some("The push would change a protected branch.".to_owned())).await;
1053 return Ok(response);
1054 }
1055 git_http::Push::Blocked(response) => {
1056 self.finish_git(audit, 403, Some("The push adds a secret.".to_owned())).await;
1057 return Ok(response);
1058 }
1059 };
1060 self.finish_git(audit, forwarded.response.status_code(), None).await;
1061
1062 // Artifacts' own push notifications are per repository, which does
1063 // not fit a repo per pull request, so the front end reports pushes
1064 // itself: one event for each branch that moved.
1065 let accepted = forwarded.response.status_code() == 200 && !forwarded.pushed.is_empty();
1066 if accepted && let Some(repo) = self.registry.by_path(&git.path).await? {
1067 let stored = self.store.open(&store_key(&repo)).await?;
1068 let actor = viewer.map(|user: User| user.id);
1069 for pushed in &forwarded.pushed {
1070 // The store can refuse one ref and accept another, so each
1071 // branch is checked against where it actually is. A tag the
1072 // store cannot read back is taken as pushed.
1073 let moved = match pushed.branch() {
1074 Some(branch) => stored
1075 .log(branch, 1)
1076 .await?
1077 .first()
1078 .is_some_and(|commit| commit.hash == pushed.after),
1079 None => stored.log(&pushed.git_ref, 1).await.map_or(true, |head| {
1080 head.first().is_none_or(|commit| commit.hash == pushed.after)
1081 }),
1082 };
1083 if moved {
1084 self.publish_push(
1085 &repo,
1086 &pushed.git_ref,
1087 pushed.before.as_deref(),
1088 &pushed.after,
1089 actor.clone(),
1090 )
1091 .await?;
1092 }
1093 }
1094 }
1095 Ok(forwarded.response)
1096 }
1097}
1098
1099#[event(fetch)]
1100async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
1101 let repos = Repos {
1102 registry: Registry { db: env.d1("DB")? },
1103 store: ArtifactsStore::new(&env)?,
1104 events: env.service("EVENTS")?,
1105 security: env.service("SECURITY").ok(),
1106 };
1107 let Some(method) = rpc_method(&request) else {
1108 return repos.git_http(request, &env).await;
1109 };
1110 let body: serde_json::Value = request.json().await?;
1111
1112 match method.as_str() {
1113 "get" => reply(&repos.get(args(body)?).await?),
1114 "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
1115 "readable" => {
1116 let a: ReadableArgs = args(body)?;
1117 reply(&repos.registry.readable(&a.ids, &a.viewer).await?)
1118 }
1119 "public_namespaces" => {
1120 let a: PublicNamespacesArgs = args(body)?;
1121 reply(&repos.registry.public_namespaces(&a.owner_id).await?)
1122 }
1123 "path_by_id" => {
1124 let a: PathByIdArgs = args(body)?;
1125 reply(
1126 &repos
1127 .registry
1128 .by_id(&a.id)
1129 .await?
1130 .filter(|repo| repo.fork_of.is_none())
1131 .map(|repo| RepoPath {
1132 namespace: repo.namespace,
1133 name: repo.name,
1134 }),
1135 )
1136 }
1137 "list" => {
1138 let a: ListArgs = args(body)?;
1139 reply(
1140 &repos
1141 .registry
1142 .list(
1143 &a.viewer,
1144 a.query.as_deref(),
1145 a.namespace.as_deref(),
1146 a.member_only,
1147 )
1148 .await?,
1149 )
1150 }
1151 "create" => reply(&repos.create(args(body)?).await?),
1152 "update" => reply(&repos.update(args(body)?).await?),
1153 "tree" => reply(&repos.tree(args(body)?).await?),
1154 "blob" => reply(&repos.blob(args(body)?).await?),
1155 "log" => reply(&repos.log(args(body)?).await?),
1156 "blame" => reply(&repos.blame(args(body)?).await?),
1157 "fork_for_pull" => reply(&repos.fork_for_pull(args(body)?).await?),
1158 "git_access" => reply(&repos.git_access(args(body)?).await?),
1159 "branches" => reply(&repos.branches(args(body)?).await?),
1160 "head" => reply(&repos.head(args(body)?).await?),
1161 "behind" => reply(&repos.behind(args(body)?).await?),
1162 "divergence" => reply(&repos.divergence(args(body)?).await?),
1163 "land" => reply(&repos.land(args(body)?).await?),
1164 "delete_branch" => reply(&repos.delete_branch(args(body)?).await?),
1165 "compare" => reply(&repos.compare(args(body)?).await?),
1166 "scan_history" => reply(&repos.scan_history(args(body)?).await?),
1167 "find_lockfiles" => reply(&repos.find_lockfiles(args(body)?).await?),
1168 "list_files" => reply(&repos.list_files(args(body)?).await?),
1169 "changed_files" => reply(&repos.changed_files(args(body)?).await?),
1170 "read_blobs" => reply(&repos.read_blobs(args(body)?).await?),
1171 "all_ids" => {
1172 let a: AllIdsArgs = args(body)?;
1173 let limit = a.limit.clamp(1, 500);
1174 let ids = repos.registry.ids_after(a.after.as_deref(), limit).await?;
1175 let next = (ids.len() == limit as usize).then(|| ids.last().cloned()).flatten();
1176 reply(&IdPage { ids, next })
1177 }
1178 _ => Response::error("Unknown method", 404),
1179 }
1180}
1181
1182/// Events from the bus. Only a workspace's rename concerns this service:
1183/// its repositories move to the workspace's current slug, asked of identity
1184/// by id, so a repeated or late delivery lands in the same place. Their git
1185/// store keys stay as they were.
1186#[event(queue)]
1187async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
1188 let registry = Registry { db: env.d1("DB")? };
1189 let identity = env.service("IDENTITY")?;
1190 for message in batch.messages()? {
1191 let event = message.body();
1192 if event.kind != "workspace.renamed" {
1193 continue;
1194 }
1195 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
1196 worker::console_error!("workspace.renamed {} could not be read", event.id);
1197 continue;
1198 };
1199 let names: HashMap<String, String> = g1t_kit::call(
1200 &identity,
1201 "usernames",
1202 &g1t_contracts::identity::UsernamesArgs {
1203 ids: vec![renamed.workspace_id.clone()],
1204 },
1205 )
1206 .await?;
1207 let current = names
1208 .get(&renamed.workspace_id)
1209 .cloned()
1210 .unwrap_or_else(|| renamed.to.clone());
1211 let left = registry
1212 .rename_namespace(&renamed.stale_slugs(&current), &current)
1213 .await?;
1214 if left > 0 {
1215 worker::console_error!(
1216 "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
1217 renamed.from,
1218 renamed.to
1219 );
1220 }
1221 }
1222 Ok(())
1223}