pr_01m47d24b0e6n91zwymwxg0vpx/services/repos/src/lib.rs

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