g1t/services/repos/src/lib.rs

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