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