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