g1t/services/repos/src/lib.rs

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