Skip to content

g1t/services/repos/src/lib.rs

2,621 lines113,000 bytesCodeBlameRaw
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 about;
9mod backups;
10mod blame;
11mod catch_up;
12mod coalesce;
13mod commit_file;
14mod contributors;
15mod diff;
16mod fallback;
17mod forks;
18mod git_http;
19mod git_ops;
20mod import;
21mod land;
22mod languages;
23mod last_commits;
24mod license;
25mod lifecycle;
26mod listing;
27mod meters;
28mod mirror;
29mod moves;
30mod namespaces;
31mod pack_cache;
32mod pack_limits;
33mod refs;
34mod refs_cache;
35mod registry;
36mod resilience;
37mod rule_facts;
38mod rules;
39mod run_access;
40mod secret_scan;
41mod shards;
42mod shared;
43mod signatures;
44mod stats;
45mod store;
46mod transfer;
47
48use g1t_contracts::events::{
49 Event, GitPush, NewEvent, Publish, RepoCreated, RepoForked, RepoUpdated, WorkspaceDeleted,
50 WorkspaceDeleting, WorkspaceRenamed, WorkspaceRestored,
51};
52use g1t_contracts::access::{self, Capability};
53use g1t_contracts::repos::*;
54use g1t_contracts::time::rfc3339;
55use g1t_contracts::{FailureCode, Outcome, PrincipalKind, User, Viewer, is_valid_repo_name, new_id};
56use g1t_kit::{args, now_ms, reply, rpc_method};
57use std::collections::{HashMap, HashSet, VecDeque};
58use std::rc::Rc;
59
60use serde::Serialize;
61use worker::{
62 Context, Env, Fetcher, MessageBatch, Method, Request, Response, Result, ScheduleContext, ScheduledEvent,
63 event,
64};
65
66use registry::{Registry, can_read, can_write, store_key};
67use store::{ArtifactsStore, GitRepo, GitStore, Scope};
68
69/// Namespace that holds every pull request's fork: `pulls/<pull id>`.
70pub(crate) const PULLS_NAMESPACE: &str = "pulls";
71const MAX_TEXT_BYTES: usize = 512 * 1024;
72/// How far back a pull request may have forked and still be landed.
73const MAX_ANCESTRY: u32 = 1000;
74/// The most tags a repository's Tags page reads and lists.
75const MAX_TAGS_READ: usize = 100;
76
77/// One path segment, percent-encoded for a cache key.
78fn urlencoding_segment(segment: &str) -> String {
79 segment
80 .bytes()
81 .map(|b| if b.is_ascii_alphanumeric() || b"-._~".contains(&b) { (b as char).to_string() } else { format!("%{b:02X}") })
82 .collect()
83}
84const MAX_DESCRIPTION_CHARS: usize = 200;
85pub(crate) const SOURCE: &str = "repos";
86pub(crate) const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
87
88pub(crate) fn not_found<T>() -> Outcome<T> {
89 Outcome::fail(FailureCode::NotFound, "Repository not found.")
90}
91
92/// Decoded text, or `None` when the file is too large or looks binary.
93/// Whether a ref is a full commit hash rather than a branch name.
94fn is_commit_hash(git_ref: &str) -> bool {
95 git_ref.len() == 40 && git_ref.bytes().all(|b| b.is_ascii_hexdigit())
96}
97
98fn text_of(bytes: Vec<u8>) -> Option<String> {
99 if bytes.len() > MAX_TEXT_BYTES || bytes.contains(&0) {
100 return None;
101 }
102 Some(String::from_utf8_lossy(&bytes).into_owned())
103}
104
105fn is_readme(name: &str) -> bool {
106 matches!(
107 name.to_lowercase().as_str(),
108 "readme" | "readme.md" | "readme.markdown" | "readme.txt"
109 )
110}
111
112/// Whether `ancestor` is reachable from the newest commit in `history`.
113///
114/// `history` is the first-parent chain, which is all the store lists; a fork
115/// that merged the target branch in has the target's head on a second
116/// parent, so the walk follows every parent.
117async fn descends_from<R: GitRepo>(repo: &R, history: &[Commit], ancestor: &str) -> Result<bool> {
118 let known: HashMap<&str, &[String]> = history
119 .iter()
120 .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
121 .collect();
122 let mut seen = HashSet::new();
123 let mut queue: Vec<String> = history
124 .first()
125 .map(|c| c.hash.clone())
126 .into_iter()
127 .collect();
128 while let Some(hash) = queue.pop() {
129 if hash == ancestor {
130 return Ok(true);
131 }
132 if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
133 continue;
134 }
135 match known.get(hash.as_str()) {
136 Some(parents) => queue.extend(parents.iter().cloned()),
137 None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
138 }
139 }
140 Ok(false)
141}
142
143/// The commit closest to the newest in `history` that is also in `shared`:
144/// where a fork and the repository it came from last agreed.
145async fn nearest_ancestor_in<R: GitRepo>(
146 repo: &R,
147 history: &[Commit],
148 shared: &HashSet<String>,
149) -> Result<Option<String>> {
150 let known: HashMap<&str, &[String]> = history
151 .iter()
152 .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
153 .collect();
154 let mut seen = HashSet::new();
155 let mut queue: VecDeque<String> = history
156 .first()
157 .map(|c| c.hash.clone())
158 .into_iter()
159 .collect();
160 while let Some(hash) = queue.pop_front() {
161 if shared.contains(&hash) {
162 return Ok(Some(hash));
163 }
164 if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
165 continue;
166 }
167 match known.get(hash.as_str()) {
168 Some(parents) => queue.extend(parents.iter().cloned()),
169 None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
170 }
171 }
172 Ok(None)
173}
174
175thread_local! {
176 /// Targets' sides of mergeability, by head (coalesce.rs).
177 static TARGETS: std::cell::RefCell<coalesce::Memo<coalesce::TargetKey, Rc<coalesce::TargetSide>>> =
178 std::cell::RefCell::new(coalesce::Memo::new(coalesce::TARGET_TTL_MS, 32));
179 /// What targets changed between two trees.
180 static THEIRS: std::cell::RefCell<coalesce::Memo<coalesce::TheirsKey, (Vec<String>, bool)>> =
181 std::cell::RefCell::new(coalesce::Memo::new(coalesce::THEIRS_TTL_MS, 256));
182 /// What repositories hold, as read for a push's first request, for the
183 /// same push's second: a push's POST does not wait on the database.
184 static HELD: std::cell::RefCell<coalesce::Memo<String, u64>> =
185 std::cell::RefCell::new(coalesce::Memo::new(60_000, 512));
186}
187
188pub(crate) struct Repos<S: GitStore> {
189 registry: Registry,
190 store: S,
191 events: Fetcher,
192 /// Asked during a push which secrets have been allowed.
193 security: Option<Fetcher>,
194 /// Asked whether a workspace is on a plan, for its private storage.
195 billing: Option<Fetcher>,
196 /// Says which rulesets hold for a change to a branch or tag, and keeps
197 /// how they judged it (rules.rs). `None` where it is not deployed: the
198 /// old protection flag then holds on push.
199 work: Option<Fetcher>,
200 /// Told when a repository moves, for the tokens of agents at work on it.
201 identity: Option<Fetcher>,
202 /// What a free workspace's private repositories may hold.
203 free_private_bytes: i64,
204 /// Days a pull request's working copy is kept after it settles (forks.rs).
205 pub(crate) fork_days: u64,
206 /// The most a repository may hold (pack_limits.rs), and what happens
207 /// to a push too large to scan.
208 repo_limit: u64,
209 large_pushes: git_http::LargePushes,
210 /// Which git store namespace new repositories go in (shards.rs), and
211 /// the most each should hold (`ARTIFACTS_NAMESPACE_LIMITS`).
212 placement: shards::Placement,
213 limits: HashMap<String, u64>,
214 /// What isolates share: answers that list refs (refs_cache.rs).
215 shared: Option<Rc<shared::Shared>>,
216 /// Packs for fresh clones (pack_cache.rs); `None` without the bucket.
217 packs: Option<Rc<pack_cache::Packs>>,
218}
219
220impl<S: GitStore> Repos<S> {
221 /// Records that the refs of the repository with this id changed, once
222 /// they have, so that the answers kept that list them go stale (see
223 /// refs_cache.rs). Everything that changes a repository's refs calls
224 /// this after it (`every_ref_writer_records_the_change` checks). A
225 /// failure is logged: the change itself happened, and what was kept
226 /// expires within `refs_cache::TTL_SECONDS` regardless.
227 pub(crate) async fn refs_moved(&self, repo_id: &str) {
228 if let Err(error) = self.registry.refs_moved(repo_id).await {
229 worker::console_error!("refs of {repo_id} changed but not recorded: {error}");
230 }
231 }
232
233 pub(crate) async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> {
234 g1t_kit::call(
235 &self.events,
236 "publish",
237 &Publish {
238 events: vec![event],
239 },
240 )
241 .await
242 }
243
244 /// Whether the viewer may read `repo`. A pull request's fork of a
245 /// private repository can be read by everyone who can read that
246 /// repository, so its members can review and check out the change, as
247 /// well as by whoever opened the pull request.
248 async fn may_read(&self, repo: &Repo, viewer: &Viewer) -> Result<bool> {
249 if can_read(repo, viewer) {
250 return Ok(true);
251 }
252 let Some(source_id) = &repo.fork_of else {
253 return Ok(false);
254 };
255 Ok(self
256 .registry
257 .by_id(source_id)
258 .await?
259 .is_some_and(|source| can_read(&source, viewer)))
260 }
261
262 /// `repo`, if there is one and the viewer may read it.
263 async fn visible(&self, repo: Option<Repo>, viewer: &Viewer) -> Result<Option<Repo>> {
264 Ok(match repo {
265 Some(repo) if self.may_read(&repo, viewer).await? => Some(repo),
266 _ => None,
267 })
268 }
269
270 /// Resolves a repo the viewer may read; private repos look missing.
271 pub(crate) async fn readable(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
272 self.visible(self.registry.by_path(path).await?, viewer)
273 .await
274 }
275
276 async fn get(&self, a: GetArgs) -> Result<Outcome<Repo>> {
277 Ok(self
278 .readable(&a.path, &a.viewer)
279 .await?
280 .map_or_else(not_found, Outcome::Ok))
281 }
282
283 async fn get_by_id(&self, a: GetByIdArgs) -> Result<Outcome<Repo>> {
284 Ok(self
285 .visible(self.registry.by_id(&a.id).await?, &a.viewer)
286 .await?
287 .map_or_else(not_found, Outcome::Ok))
288 }
289
290 async fn update(&self, a: UpdateArgs) -> Result<Outcome<Repo>> {
291 let viewer = Some(a.actor.clone());
292 let Some(repo) = self.readable(&a.path, &viewer).await? else {
293 return Ok(not_found());
294 };
295 // Its details take Maintain; its protection, Maintain too; who can
296 // see it, Admin (below). See g1t_contracts::access.
297 let protection_changes = a.protected.is_some_and(|protected| protected != repo.protected);
298 let details_change = a.description.is_some() || a.website.is_some() || a.topics.is_some();
299 let mut needed = Vec::new();
300 if details_change || !protection_changes {
301 needed.push(Capability::ManageSettings);
302 }
303 if protection_changes {
304 needed.push(Capability::ManageProtection);
305 }
306 let full_name = format!("{}/{}", repo.namespace, repo.name);
307 if repo.fork_of.is_some() {
308 return Ok(Outcome::fail(FailureCode::Forbidden, access::needs(Capability::ManageSettings, &full_name)));
309 }
310 if let Some(missing) = needed.into_iter().find(|capability| !registry::can(&repo, &viewer, *capability)) {
311 return Ok(Outcome::fail(FailureCode::Forbidden, access::needs(missing, &full_name)));
312 }
313 if !a.actor.verified {
314 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
315 }
316 if let Some((code, message)) = lifecycle::archived_refusal(&repo) {
317 return Ok(Outcome::fail(code, message));
318 }
319 let description = match a.description {
320 Some(text) => Some(
321 text.trim()
322 .chars()
323 .take(MAX_DESCRIPTION_CHARS)
324 .collect::<String>(),
325 )
326 .filter(|text| !text.is_empty()),
327 None => repo.description.clone(),
328 };
329 let website = match a.website.as_deref() {
330 Some(text) => match clean_website(text) {
331 Ok(website) => website,
332 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
333 },
334 None => repo.website.clone(),
335 };
336 // Who can see it is an owner's to change, and a free workspace's
337 // storage may not take it private: see lifecycle.rs.
338 let wants_private = a.is_private.filter(|private| *private != repo.is_private);
339 if wants_private.is_some()
340 && let Err((code, message)) = lifecycle::admin_only(
341 lifecycle::Asker::on(&a.actor, &repo),
342 &repo.namespace,
343 "change the visibility of",
344 Capability::Administer,
345 )
346 {
347 return Ok(Outcome::fail(code, message));
348 }
349 let is_private = repo.is_private;
350 let protected = a.protected.unwrap_or(repo.protected);
351 let topics = match &a.topics {
352 Some(topics) => match clean_topics(topics) {
353 Ok(topics) => topics,
354 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
355 },
356 None => repo.topics.clone(),
357 };
358 self.registry
359 .update(&repo.id, description.as_deref(), protected, &topics, website.as_deref())
360 .await?;
361 let updated = Repo {
362 description,
363 is_private,
364 protected,
365 topics,
366 website,
367 ..repo
368 };
369 // Whether the default branch takes only pull requests is now its
370 // branch protection ruleset's to say (work's rulesets.rs).
371 if let (Some(protected), Some(work)) = (a.protected, &self.work) {
372 #[derive(Serialize)]
373 struct RequirePullRequest<'a> {
374 repo: &'a Repo,
375 protected: bool,
376 actor: &'a User,
377 }
378 let set: Result<Outcome<bool>> =
379 g1t_kit::call(work, "set_requires_pull_request", &RequirePullRequest { repo: &updated, protected, actor: &a.actor }).await;
380 match set {
381 Ok(Outcome::Ok(_)) => {}
382 Ok(Outcome::Fail(failure)) => return Ok(Outcome::Fail(failure)),
383 Err(error) => return Err(error),
384 }
385 }
386 if let Some(private) = wants_private {
387 return self.change_visibility(updated, private, &a.actor, a.surface).await;
388 }
389 let visibility_changed = false;
390 // Search and anything else that shows the repository hears of it;
391 // a change of visibility is announced on its own as well, so that
392 // what was public stops being shown at once.
393 self.publish(NewEvent {
394 kind: "repo.updated",
395 source: SOURCE,
396 repo_id: Some(updated.id.clone()),
397 actor: Some(a.actor.id.clone()),
398 data: RepoUpdated {
399 repo_id: updated.id.clone(),
400 namespace: updated.namespace.clone(),
401 name: updated.name.clone(),
402 is_private,
403 visibility_changed,
404 },
405 })
406 .await?;
407 Ok(Outcome::Ok(updated))
408 }
409
410 /// The repository with this id, if it is not a fork, and its store.
411 async fn stored(&self, repo_id: &str) -> Result<Option<S::Repo>> {
412 match self.registry.by_id(repo_id).await? {
413 Some(repo) if repo.fork_of.is_none() => Ok(Some(self.store.open(&store_key(&repo)).await?)),
414 _ => Ok(None),
415 }
416 }
417
418 async fn list_files(&self, a: ListFilesArgs) -> Result<FileList> {
419 let Some(repo) = self.registry.by_id(&a.repo_id).await?.filter(|repo| repo.fork_of.is_none()) else {
420 return Ok(FileList::default());
421 };
422 let git = self.store.open(&store_key(&repo)).await?;
423 let head = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
424 listing::list(&git, None, &head, &a.skip_dirs, a.limit).await
425 }
426
427 async fn changed_files(&self, a: ChangedFilesArgs) -> Result<FileList> {
428 let Some(git) = self.stored(&a.repo_id).await? else {
429 return Ok(FileList::default());
430 };
431 listing::list(&git, a.base.as_deref(), &a.head, &a.skip_dirs, a.limit).await
432 }
433
434 /// Branches and tags with their commits, for g1t's own services.
435 async fn refs_of(&self, a: RefsArgs) -> Result<Option<RepoRefs>> {
436 let Some(repo) = self.registry.by_id(&a.repo_id).await?.filter(|repo| repo.fork_of.is_none()) else {
437 return Ok(None);
438 };
439 let git = self.store.open(&store_key(&repo)).await?;
440 let access = git.access(Scope::Read).await?;
441 let refs = refs::heads_and_tags(refs::all(&access).await?)
442 .into_iter()
443 .map(|(name, commit)| GitRefEntry { name, commit })
444 .collect();
445 Ok(Some(RepoRefs { repo, refs }))
446 }
447
448 async fn raw_file(&self, a: RawFileArgs) -> Result<Option<RawFile>> {
449 use base64::Engine;
450 let Some(git) = self.stored(&a.repo_id).await? else {
451 return Ok(None);
452 };
453 Ok(git
454 .read_file(&a.git_ref, &a.path)
455 .await?
456 .filter(|bytes| bytes.len() <= a.max_bytes as usize)
457 .map(|bytes| RawFile { size: bytes.len() as u64, data: base64::engine::general_purpose::STANDARD.encode(bytes) }))
458 }
459
460 async fn raw_blobs(&self, a: RawBlobsArgs) -> Result<Vec<RawBlob>> {
461 use base64::Engine;
462 let Some(git) = self.stored(&a.repo_id).await? else {
463 return Ok(Vec::new());
464 };
465 let hashes: Vec<&String> = a.hashes.iter().take(MAX_READ_BLOBS).collect();
466 let mut out = Vec::with_capacity(hashes.len());
467 // A few at a time, as listing::read does: each is a round trip.
468 for group in hashes.chunks(8) {
469 let read = futures_util::future::try_join_all(group.iter().map(|hash| git.read_blob(hash))).await?;
470 for (hash, bytes) in group.iter().zip(read) {
471 let size = bytes.as_ref().map_or(0, |bytes| bytes.len() as u64);
472 let data = bytes
473 .filter(|bytes| bytes.len() <= a.max_bytes as usize)
474 .map(|bytes| base64::engine::general_purpose::STANDARD.encode(bytes));
475 out.push(RawBlob { hash: (*hash).clone(), size, data });
476 }
477 }
478 Ok(out)
479 }
480
481 async fn read_blobs(&self, a: ReadBlobsArgs) -> Result<Vec<BlobText>> {
482 let Some(git) = self.stored(&a.repo_id).await? else {
483 return Ok(Vec::new());
484 };
485 listing::read(&git, &a.hashes, a.max_bytes.min(MAX_TEXT_BYTES as u32)).await
486 }
487
488 async fn create(&self, a: CreateArgs) -> Result<Outcome<Repo>> {
489 if !a.owner.verified {
490 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
491 }
492 let name = a.name.trim().to_lowercase();
493 if !is_valid_repo_name(&name) {
494 return Ok(Outcome::fail(
495 FailureCode::Invalid,
496 "Use letters, digits, dots, hyphens and underscores only.",
497 ));
498 }
499 let namespace = a.namespace.trim().to_lowercase();
500 if namespace.is_empty() {
501 return Ok(Outcome::fail(
502 FailureCode::Invalid,
503 "Say which workspace to create the repository in.",
504 ));
505 }
506 if !a.owner.is_member(&namespace) {
507 return Ok(Outcome::fail(
508 FailureCode::Forbidden,
509 "You are not a member of that workspace.",
510 ));
511 }
512 let path = RepoPath { namespace, name };
513 match self.registry.by_path_any(&path).await? {
514 Some((_, None)) => {
515 return Ok(Outcome::fail(
516 FailureCode::Conflict,
517 "That workspace already has a repository with that name.",
518 ));
519 }
520 Some((_, Some(_))) => {
521 return Ok(Outcome::fail(
522 FailureCode::Conflict,
523 format!(
524 "{}/{} 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.",
525 path.namespace, path.name
526 ),
527 ));
528 }
529 None => {}
530 }
531 // With a credential (a GitHub App installation's token), everything
532 // is copied: every branch and tag. See mirror.rs.
533 let mut credentialed = None;
534 if let (Some(url), Some(token)) = (a.import_url.as_deref(), a.import_token.as_deref()) {
535 let Some(url) = import::clean_url(url) else {
536 return Ok(Outcome::fail(FailureCode::Invalid, "That is not an https repository address."));
537 };
538 let source = mirror::Endpoint::github(&url, token);
539 match mirror::probe(&source).await? {
540 Ok(advertised) => credentialed = Some((source, advertised)),
541 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
542 }
543 }
544 // An import is fetched before anything is created, so that an
545 // address that does not work leaves nothing behind.
546 let mut imported = None;
547 if let Some(url) = a
548 .import_url
549 .as_deref()
550 .map(str::trim)
551 .filter(|url| !url.is_empty() && credentialed.is_none())
552 {
553 let Some(url) = import::clean_url(url) else {
554 return Ok(Outcome::fail(
555 FailureCode::Invalid,
556 "Give the https address of a public repository, such as https://github.com/owner/repo.",
557 ));
558 };
559 let remote = match import::discover(&url).await? {
560 Ok(remote) => remote,
561 Err(reason) => return Ok(Outcome::fail(FailureCode::Invalid, reason)),
562 };
563 imported = Some((remote, url));
564 }
565 let now = now_ms();
566 let repo = Repo {
567 id: new_id("rep", now),
568 namespace: path.namespace,
569 name: path.name,
570 description: a
571 .description
572 .map(|text| text.trim().to_owned())
573 .filter(|text| !text.is_empty()),
574 is_private: a.is_private,
575 owner_id: a.owner.id.clone(),
576 default_branch: imported
577 .as_ref()
578 .map(|(remote, _)| remote.branch.clone())
579 .or_else(|| credentialed.as_ref().and_then(|(_, advertised)| advertised.default_branch()))
580 .unwrap_or_else(|| "main".to_owned()),
581 fork_of: None,
582 protected: false,
583 created_at: rfc3339(now),
584 topics: Vec::new(),
585 website: None,
586 archived_at: None,
587 };
588 let namespace = match self.place(&repo).await? {
589 Ok(namespace) => namespace,
590 Err(unplaced) => return Ok(Outcome::fail(FailureCode::Conflict, unplaced.message())),
591 };
592 self.registry
593 .claim_store_key(&repo, namespace.as_deref(), &self.store.default_namespace())
594 .await?;
595 self.store
596 .create(
597 &store_key(&repo),
598 repo.description.as_deref(),
599 &repo.default_branch,
600 )
601 .await?;
602 self.registry.insert(&repo).await?;
603 // A repository that was transferred away from this path stops
604 // redirecting here.
605 self.registry
606 .drop_redirect(&RepoPath {
607 namespace: repo.namespace.clone(),
608 name: repo.name.clone(),
609 })
610 .await?;
611 // Every branch and tag the import made, announced as pushes.
612 let mut pushed: Vec<(String, String)> = Vec::new();
613 // A public repository, read with no credential: every branch and
614 // tag is copied too, the default branch the one its HEAD names.
615 if let Some((_, url)) = imported {
616 let access = self
617 .store
618 .open(&store_key(&repo))
619 .await?
620 .access(Scope::Write)
621 .await?;
622 let target = mirror::Endpoint::bearer(&access.remote, &access.token);
623 let copied = mirror::copy(&mirror::Endpoint::anonymous(&url), &target, mirror::Prune::Yes).await?;
624 self.refs_moved(&repo.id).await;
625 match copied {
626 Ok(copied) => pushed = mirror::import_pushes(&copied.updated, &repo.default_branch),
627 Err(reason) => {
628 self.registry.remove(&repo.id).await?;
629 return Ok(Outcome::fail(
630 FailureCode::Invalid,
631 format!("The repository could not be stored: {reason}"),
632 ));
633 }
634 }
635 }
636 if let Some((source, _)) = credentialed {
637 let access = self
638 .store
639 .open(&store_key(&repo))
640 .await?
641 .access(Scope::Write)
642 .await?;
643 let target = mirror::Endpoint::bearer(&access.remote, &access.token);
644 let copied = mirror::copy(&source, &target, mirror::Prune::Yes).await?;
645 self.refs_moved(&repo.id).await;
646 match copied {
647 Ok(copied) => pushed = mirror::import_pushes(&copied.updated, &repo.default_branch),
648 Err(reason) => {
649 self.registry.remove(&repo.id).await?;
650 return Ok(Outcome::fail(
651 FailureCode::Invalid,
652 format!("The repository could not be copied: {reason}"),
653 ));
654 }
655 }
656 }
657 self.publish(NewEvent {
658 kind: "repo.created",
659 source: SOURCE,
660 repo_id: Some(repo.id.clone()),
661 actor: Some(a.owner.id),
662 data: RepoCreated {
663 repo_id: repo.id.clone(),
664 namespace: repo.namespace.clone(),
665 name: repo.name.clone(),
666 is_private: repo.is_private,
667 },
668 })
669 .await?;
670 for (git_ref, head) in &pushed {
671 self.publish_push(&repo, git_ref, None, head, None).await?;
672 }
673 Ok(Outcome::Ok(repo))
674 }
675
676 /// Where a workspace keeps its data, asked of identity only when an EU
677 /// namespace is configured: without one, every workspace's
678 /// repositories go anywhere and identity is never asked.
679 async fn residency_of(&self, workspace: &str) -> Result<shards::Residency> {
680 if self.placement.eu.is_none() {
681 return Ok(shards::Residency::Anywhere);
682 }
683 let Some(identity) = &self.identity else {
684 return Ok(shards::Residency::Anywhere);
685 };
686 let residency: Option<g1t_contracts::identity::DataResidency> = g1t_kit::call(
687 identity,
688 "workspace_residency",
689 &g1t_contracts::identity::SlugArgs { slug: workspace.to_owned() },
690 )
691 .await?;
692 Ok(match residency {
693 Some(g1t_contracts::identity::DataResidency::Eu) => shards::Residency::Eu,
694 _ => shards::Residency::Anywhere,
695 })
696 }
697
698 /// How each bound namespace stands (namespaces.rs).
699 async fn standings(&self) -> Result<Vec<namespaces::Standing>> {
700 let bound = self.store.namespaces();
701 let default = self.store.default_namespace();
702 let now = now_ms();
703 let config = namespaces::Configured {
704 bound: &bound,
705 default: &default,
706 placement: &self.placement,
707 limits: &self.limits,
708 on_fallback: &|namespace| self.store.on_fallback(&shards::compose(Some(namespace), "x", &default)),
709 writable: &|namespace| self.store.writable(namespace),
710 breaker_open: &|namespace| resilience::open_now(namespace, now),
711 };
712 let (held, recent) = futures_util::future::join(namespaces::held(&self.registry.db), namespaces::recent(&self.registry.db, now)).await;
713 Ok(namespaces::standings(&config, &held?, &recent?))
714 }
715
716 /// The namespace a new repository goes in (shards.rs): its workspace's
717 /// residency, then how each namespace stands, read only when there is
718 /// a choice to make. `Ok(None)` for the default.
719 async fn place(&self, repo: &Repo) -> Result<std::result::Result<Option<String>, shards::Unplaced>> {
720 let residency = self.residency_of(&repo.namespace).await?;
721 let bound = self.store.namespaces();
722 let loads = if residency == shards::Residency::Anywhere && !self.placement.needs_loads(&bound) {
723 // One namespace to choose from at most: nothing to read.
724 bound
725 .iter()
726 .map(|namespace| shards::Load {
727 namespace: namespace.clone(),
728 bound: true,
729 writable: self.store.writable(namespace),
730 ..shards::Load::default()
731 })
732 .collect()
733 } else {
734 let default = self.store.default_namespace();
735 let now = now_ms();
736 let config = namespaces::Configured {
737 bound: &bound,
738 default: &default,
739 placement: &self.placement,
740 limits: &self.limits,
741 on_fallback: &|namespace| self.store.on_fallback(&shards::compose(Some(namespace), "x", &default)),
742 writable: &|namespace| self.store.writable(namespace),
743 breaker_open: &|namespace| resilience::open_now(namespace, now),
744 };
745 namespaces::loads(&self.registry.db, &config, now).await?
746 };
747 Ok(self.placement.choose(&repo.id, residency, &loads))
748 }
749
750 /// `storage_options`: what a workspace may choose about where its
751 /// repositories are kept.
752 fn storage_options(&self) -> StorageOptions {
753 let bound = self.store.namespaces();
754 StorageOptions { eu_available: self.placement.eu_available(&bound, |namespace| self.store.writable(namespace)) }
755 }
756
757 async fn tree(&self, a: TreeArgs) -> Result<Outcome<TreeView>> {
758 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
759 return Ok(not_found());
760 };
761 let git = self.read_git(&repo).await?;
762 let git_ref = a
763 .git_ref
764 .clone()
765 .unwrap_or_else(|| repo.default_branch.clone());
766
767 let Some(head) = git.log(&git_ref, 1).await?.into_iter().next() else {
768 // An unknown ref is an error; a repo with no commits is just empty.
769 if a.git_ref.is_some() {
770 return Ok(Outcome::fail(
771 FailureCode::NotFound,
772 "No such branch, tag or commit.",
773 ));
774 }
775 return Ok(Outcome::Ok(TreeView {
776 repo,
777 git_ref,
778 path: a.tree_path,
779 head: None,
780 entries: Vec::new(),
781 readme: None,
782 }));
783 };
784
785 let no_directory = || Outcome::fail(FailureCode::NotFound, "No such directory.");
786 let mut entries = git.read_tree(&head.tree_hash).await?;
787 for segment in a.tree_path.split('/').filter(|segment| !segment.is_empty()) {
788 let next = entries.as_ref().and_then(|entries| {
789 entries
790 .iter()
791 .find(|entry| entry.name == segment && entry.kind == EntryKind::Tree)
792 });
793 let Some(next) = next else {
794 return Ok(no_directory());
795 };
796 entries = git.read_tree(&next.hash).await?;
797 }
798 let Some(mut entries) = entries else {
799 return Ok(no_directory());
800 };
801 // Directories first, then by name.
802 entries.sort_by(|a, b| {
803 (b.kind == EntryKind::Tree)
804 .cmp(&(a.kind == EntryKind::Tree))
805 .then_with(|| a.name.cmp(&b.name))
806 });
807
808 let readme_entry = entries
809 .iter()
810 .find(|entry| entry.kind == EntryKind::Blob && is_readme(&entry.name));
811 let readme = match readme_entry {
812 Some(entry) => git.read_blob(&entry.hash).await?.map(|bytes| Readme {
813 name: entry.name.clone(),
814 text: text_of(bytes),
815 }),
816 None => None,
817 };
818 Ok(Outcome::Ok(TreeView {
819 repo,
820 git_ref,
821 path: a.tree_path,
822 head: Some(head),
823 entries,
824 readme,
825 }))
826 }
827
828 async fn blob(&self, a: BlobArgs) -> Result<Outcome<BlobView>> {
829 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
830 return Ok(not_found());
831 };
832 let bytes = if a.file_path.is_empty() {
833 None
834 } else {
835 let git = self.read_git(&repo).await?;
836 git.read_file(&a.git_ref, &a.file_path).await?
837 };
838 let Some(bytes) = bytes else {
839 return Ok(Outcome::fail(FailureCode::NotFound, "No such file."));
840 };
841 Ok(Outcome::Ok(BlobView {
842 repo,
843 git_ref: a.git_ref,
844 path: a.file_path,
845 size: bytes.len() as u64,
846 text: text_of(bytes),
847 }))
848 }
849
850 async fn blame(&self, a: BlameArgs) -> Result<Outcome<Blame>> {
851 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
852 return Ok(not_found());
853 };
854 let git = self.read_git(&repo).await?;
855 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
856 Ok(match blame::blame(&git, &git_ref, &a.file_path).await? {
857 Some(blame) => Outcome::Ok(blame),
858 None => not_found(),
859 })
860 }
861
862 async fn log(&self, a: LogArgs) -> Result<Outcome<Vec<Commit>>> {
863 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
864 return Ok(not_found());
865 };
866 let git = self.read_git(&repo).await?;
867 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
868 Ok(Outcome::Ok(git.log(&git_ref, a.limit).await?))
869 }
870
871 /// Which commit last changed each entry of a directory. Kept in this
872 /// colo's cache by repository, head commit and path: a commit's history
873 /// never changes, so an answer is good for as long as it is kept.
874 async fn last_commits(&self, a: g1t_contracts::repos::LastCommitsArgs) -> Result<Outcome<g1t_contracts::repos::LastCommits>> {
875 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
876 return Ok(not_found());
877 };
878 let git = self.read_git(&repo).await?;
879 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
880 let Some(head) = git.log(&git_ref, 1).await?.into_iter().next() else {
881 return Ok(Outcome::fail(FailureCode::NotFound, "No such branch, tag or commit."));
882 };
883 let key = format!(
884 "https://last-commits.g1t.internal/{}/{}/{}",
885 repo.id,
886 head.hash,
887 a.tree_path.split('/').map(urlencoding_segment).collect::<Vec<_>>().join("/")
888 );
889 let cache = worker::Cache::default();
890 if let Ok(Some(mut kept)) = cache.get(key.as_str(), false).await {
891 if let Ok(found) = kept.json::<g1t_contracts::repos::LastCommits>().await {
892 return Ok(Outcome::Ok(found));
893 }
894 }
895 // Asked with a budget: past it, what was found so far, not kept.
896 let started = worker::Date::now().as_millis();
897 let budget = a.budget_ms;
898 let out_of_time = move || budget.is_some_and(|budget| worker::Date::now().as_millis().saturating_sub(started) > budget);
899 let (entries, complete) = last_commits::last_commits(&git, &head.hash, &a.tree_path, &out_of_time).await?;
900 let stopped = out_of_time();
901 let found = g1t_contracts::repos::LastCommits { entries, complete };
902 if stopped && !found.complete {
903 return Ok(Outcome::Ok(found));
904 }
905 if let Ok(mut response) = worker::Response::from_json(&found) {
906 let _ = response.headers_mut().set("cache-control", "max-age=604800");
907 let _ = cache.put(key.as_str(), response).await;
908 }
909 Ok(Outcome::Ok(found))
910 }
911
912 /// The repository's tags, newest commit first, at most 100.
913 async fn tags(&self, a: g1t_contracts::repos::TagsArgs) -> Result<Outcome<Vec<g1t_contracts::repos::Tag>>> {
914 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
915 return Ok(not_found());
916 };
917 let git = self.store.open(&store_key(&repo)).await?;
918 let access = git.access(Scope::Read).await?;
919 let named: Vec<(String, String)> = refs::heads_and_tags(refs::all(&access).await?)
920 .into_iter()
921 .filter_map(|(name, hash)| name.strip_prefix("refs/tags/").map(|tag| (tag.to_owned(), hash)))
922 .collect();
923 let read = self.read_git(&repo).await?;
924 let commits = futures_util::future::join_all(named.iter().take(MAX_TAGS_READ).map(|(_, hash)| read.log(hash, 1))).await;
925 let mut tags: Vec<g1t_contracts::repos::Tag> = named
926 .into_iter()
927 .zip(commits.into_iter().map(|found| found.ok().and_then(|list| list.into_iter().next())).chain(std::iter::repeat(None)))
928 .map(|((name, _), commit)| g1t_contracts::repos::Tag { name, commit })
929 .collect();
930 tags.sort_by(|a, b| {
931 let at = |tag: &g1t_contracts::repos::Tag| tag.commit.as_ref().map(|c| c.authored_at.clone()).unwrap_or_default();
932 at(b).cmp(&at(a)).then_with(|| b.name.cmp(&a.name))
933 });
934 tags.truncate(MAX_TAGS_READ);
935 Ok(Outcome::Ok(tags))
936 }
937
938 /// The repository's branches, default branch first.
939 async fn branches(&self, a: BranchesArgs) -> Result<Outcome<Vec<Branch>>> {
940 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
941 return Ok(not_found());
942 };
943 let mut branches = self.read_git(&repo).await?.branches().await?;
944 branches.sort_by_key(|branch| branch.name != repo.default_branch);
945 Ok(Outcome::Ok(branches))
946 }
947
948 /// Whether a pull request's source lacks commits that the branch it
949 /// would merge into has.
950 async fn behind(&self, a: BehindArgs) -> Result<bool> {
951 let Some(source) = self.registry.by_id(&a.source_id).await? else {
952 return Ok(false);
953 };
954 let target = match &source.fork_of {
955 Some(id) => self.registry.by_id(id).await?,
956 None => Some(source.clone()),
957 };
958 let Some(target) = target else {
959 return Ok(false);
960 };
961 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
962 let target_branch = a.target_branch.unwrap_or_else(|| target.default_branch.clone());
963 let target_head = self
964 .read_git(&target)
965 .await?
966 .log(&target_branch, 1)
967 .await?
968 .into_iter()
969 .next()
970 .map(|commit| commit.hash);
971 let Some(target_head) = target_head else {
972 return Ok(false);
973 };
974 let source_git = self.read_git(&source).await?;
975 let history = source_git.log(&branch, MAX_ANCESTRY).await?;
976 if history.is_empty() {
977 return Ok(false);
978 }
979 Ok(!descends_from(&source_git, &history, &target_head).await?)
980 }
981
982 /// The files a pull request's source and the default branch it would
983 /// merge into each changed since they last agreed. Where the two lists
984 /// share no file, the merge cannot conflict; where they do, it may.
985 async fn divergence(&self, a: BehindArgs) -> Result<Option<Divergence>> {
986 let Some(source) = self.registry.by_id(&a.source_id).await? else {
987 return Ok(None);
988 };
989 let target = match &source.fork_of {
990 Some(id) => self.registry.by_id(id).await?,
991 None => Some(source.clone()),
992 };
993 let Some(target) = target else {
994 return Ok(None);
995 };
996 let branch = a.branch.unwrap_or_else(|| target.default_branch.clone());
997 let target_branch = a.target_branch.unwrap_or_else(|| target.default_branch.clone());
998 let source_git = self.read_git(&source).await?;
999 let target_git = self.read_git(&target).await?;
1000 // The target's side is the same for every pull request into it, and
1001 // worked out once per head (coalesce.rs).
1002 let (history, side) = futures_util::future::try_join(
1003 source_git.log(&branch, MAX_ANCESTRY),
1004 self.target_side(&target, &target_git, &target_branch),
1005 )
1006 .await?;
1007 let target_history = &side.history;
1008 let (Some(head), Some(base)) = (history.first(), target_history.first()) else {
1009 return Ok(None);
1010 };
1011 let behind = !descends_from(&source_git, &history, &base.hash).await?;
1012 let merge_base = nearest_ancestor_in(&source_git, &history, &side.shared).await?;
1013 let mut divergence = Divergence {
1014 head: head.hash.clone(),
1015 base: base.hash.clone(),
1016 merge_base: merge_base.clone(),
1017 behind,
1018 ..Divergence::default()
1019 };
1020 let merge_base_tree = match &merge_base {
1021 Some(hash) => target_history
1022 .iter()
1023 .find(|commit| commit.hash == *hash)
1024 .map(|commit| commit.tree_hash.clone()),
1025 None => None,
1026 };
1027 let Some(merge_base_tree) = merge_base_tree else {
1028 // No common history to compare from: say nothing is known.
1029 divergence.truncated = true;
1030 return Ok(Some(divergence));
1031 };
1032 let (ours, truncated_ours) =
1033 diff::changed_paths(&source_git, Some(&merge_base_tree), &head.tree_hash).await?;
1034 divergence.ours = ours;
1035 divergence.truncated = truncated_ours;
1036 if behind {
1037 let now = now_ms();
1038 let key = (target.id.clone(), merge_base_tree.clone(), base.tree_hash.clone());
1039 let (theirs, truncated_theirs) = match THEIRS.with(|memo| memo.borrow().get(&key, now)) {
1040 Some(kept) => kept,
1041 None => {
1042 let found = diff::changed_paths(&target_git, Some(&merge_base_tree), &base.tree_hash).await?;
1043 THEIRS.with(|memo| memo.borrow_mut().put(key, found.clone(), now));
1044 found
1045 }
1046 };
1047 divergence.theirs = theirs;
1048 divergence.truncated |= truncated_theirs;
1049 }
1050 Ok(Some(divergence))
1051 }
1052
1053 /// A target branch's history from its head, worked out once per head
1054 /// for every pull request asking about it (coalesce.rs). The head is
1055 /// read under the refs version; the history by its hash, which the
1056 /// object cache keeps for good.
1057 async fn target_side<R: GitRepo>(&self, target: &Repo, git: &R, branch: &str) -> Result<Rc<coalesce::TargetSide>> {
1058 let now = now_ms();
1059 let key = refs_cache::usable(registry::refs_state(&target.id), now)
1060 .map(|version| (target.id.clone(), branch.to_owned(), version));
1061 if let Some(key) = &key
1062 && let Some(side) = TARGETS.with(|memo| memo.borrow().get(key, now))
1063 {
1064 return Ok(side);
1065 }
1066 let history = match git.log(branch, 1).await?.first() {
1067 Some(head) => git.log(&head.hash, MAX_ANCESTRY).await?,
1068 None => Vec::new(),
1069 };
1070 let side = coalesce::TargetSide::new(history);
1071 if let Some(key) = key {
1072 TARGETS.with(|memo| memo.borrow_mut().put(key, side.clone(), now));
1073 }
1074 Ok(side)
1075 }
1076
1077 async fn head(&self, a: HeadArgs) -> Result<Option<String>> {
1078 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
1079 return Ok(None);
1080 };
1081 let branch = if a.branch.is_empty() { &repo.default_branch } else { &a.branch };
1082 let git = self.read_git(&repo).await?;
1083 Ok(git
1084 .log(branch, 1)
1085 .await?
1086 .into_iter()
1087 .next()
1088 .map(|commit| commit.hash))
1089 }
1090
1091 async fn delete_branch(&self, a: DeleteBranchArgs) -> Result<Outcome<bool>> {
1092 if !a.branch.starts_with(G1T_BRANCH_PREFIX) {
1093 return Ok(Outcome::fail(
1094 FailureCode::Forbidden,
1095 "Only branches g1t made for itself can be deleted this way.",
1096 ));
1097 }
1098 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
1099 return Ok(not_found());
1100 };
1101 let repo = match self.unpaused(repo).await? {
1102 Ok(repo) => repo,
1103 Err((code, message)) => return Ok(Outcome::fail(code, message)),
1104 };
1105 self.live(&repo).await?;
1106 let git = self.store.open(&store_key(&repo)).await?;
1107 let Some(old) = git
1108 .branches()
1109 .await?
1110 .into_iter()
1111 .find(|branch| branch.name == a.branch)
1112 .map(|branch| branch.hash)
1113 else {
1114 return Ok(Outcome::Ok(false));
1115 };
1116 let access = git.access(Scope::Write).await?;
1117 let deleted = land::delete_ref(&access, &a.branch, &old).await?;
1118 self.refs_moved(&repo.id).await;
1119 if let Err(reason) = deleted {
1120 return Ok(Outcome::fail(
1121 FailureCode::Conflict,
1122 format!("{} could not be deleted: {reason}", a.branch),
1123 ));
1124 }
1125 Ok(Outcome::Ok(true))
1126 }
1127
1128 async fn fork_for_pull(&self, a: ForkArgs) -> Result<Outcome<Repo>> {
1129 let viewer = Some(a.actor.clone());
1130 let Some(source) = self
1131 .registry
1132 .by_id(&a.source_id)
1133 .await?
1134 .filter(|repo| can_read(repo, &viewer))
1135 else {
1136 return Ok(not_found());
1137 };
1138 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1139 return Ok(Outcome::fail(code, message));
1140 }
1141 // Its working copy is made in its namespace: not while it moves.
1142 let source = match self.unpaused(source).await? {
1143 Ok(source) => source,
1144 Err((code, message)) => return Ok(Outcome::fail(code, message)),
1145 };
1146 let now = now_ms();
1147 let fork = Repo {
1148 id: new_id("rep", now),
1149 namespace: PULLS_NAMESPACE.to_owned(),
1150 name: a.pull_id.clone(),
1151 description: None,
1152 // A fork is exactly as visible as the repo it came from.
1153 is_private: source.is_private,
1154 owner_id: a.actor.id.clone(),
1155 default_branch: source.default_branch.clone(),
1156 fork_of: Some(source.id.clone()),
1157 protected: false,
1158 created_at: rfc3339(now),
1159 topics: Vec::new(),
1160 website: None,
1161 archived_at: None,
1162 };
1163 // Artifacts forks within a namespace: the copy goes where its
1164 // repository is.
1165 let (namespace, _) = store::locate(&store_key(&source));
1166 self.registry
1167 .claim_store_key(&fork, Some(&namespace), &self.store.default_namespace())
1168 .await?;
1169 self.store
1170 .open(&store_key(&source))
1171 .await?
1172 .fork(&store_key(&fork))
1173 .await?;
1174 self.registry.insert(&fork).await?;
1175 self.publish(NewEvent {
1176 kind: "repo.forked",
1177 source: SOURCE,
1178 repo_id: Some(source.id.clone()),
1179 actor: Some(a.actor.id),
1180 data: RepoForked {
1181 repo_id: fork.id.clone(),
1182 source_repo_id: source.id,
1183 pull_id: a.pull_id,
1184 },
1185 })
1186 .await?;
1187 Ok(Outcome::Ok(fork))
1188 }
1189
1190 async fn git_access(&self, a: GitAccessArgs) -> Result<Outcome<GitAccess>> {
1191 let found = self.registry.by_path(&a.path).await?;
1192 Ok(match self.authorize_git(&a.path, &a.viewer, a.service, found).await? {
1193 Outcome::Ok(repo) => {
1194 self.live(&repo).await?;
1195 let write = a.service == GitService::ReceivePack;
1196 if write {
1197 // A push with this credential would not pass through
1198 // here, so nothing that lists the refs is kept until it
1199 // has expired (see refs_cache.rs).
1200 let until = now_ms() + store::CREDENTIAL_LIFE_MS + 60_000;
1201 if let Err(error) = self.registry.refs_open(&repo.id, until).await {
1202 // Before the column exists nothing is kept anyway.
1203 if registry::refs_state(&repo.id).is_some() {
1204 return Err(error);
1205 }
1206 }
1207 }
1208 let scope = if write { Scope::Write } else { Scope::Read };
1209 Outcome::Ok(self.store.handout(&store_key(&repo), scope).await?)
1210 }
1211 Outcome::Fail(failure) => Outcome::Fail(failure),
1212 })
1213 }
1214
1215 /// The repository at `path` (`found`, as just read), if the viewer may
1216 /// use `service` on it: fetch from it, or push to it. A push to a path
1217 /// with nothing there makes the repository, in a workspace the pusher
1218 /// belongs to.
1219 async fn authorize_git(
1220 &self,
1221 path: &RepoPath,
1222 viewer: &Viewer,
1223 service: GitService,
1224 found: Option<Repo>,
1225 ) -> Result<Outcome<Repo>> {
1226 let mut a = GitAccessArgs {
1227 path: path.clone(),
1228 viewer: viewer.clone(),
1229 service,
1230 };
1231 let write = a.service == GitService::ReceivePack;
1232 // An access token: pushing needs code:write, reading a private
1233 // repository code:read. A public repository reads as it would for
1234 // anyone. Which repositories a token reaches is its owner's, checked
1235 // below as for anyone.
1236 if let Some(access) = a.viewer.as_ref().and_then(|user| user.token.as_deref()).cloned() {
1237 // A workflow job's token reaches its own repository only, and
1238 // the working copies of that repository's pull requests, where
1239 // their heads are.
1240 if let Some(refused) = g1t_contracts::scopes::decide_repo(&access, &format!("{}/{}", path.namespace, path.name)) {
1241 let source = match found.as_ref().and_then(|repo| repo.fork_of.as_deref()) {
1242 Some(source_id) => self.registry.by_id(source_id).await?,
1243 None => None,
1244 };
1245 if !source.is_some_and(|source| access.reaches(&format!("{}/{}", source.namespace, source.name))) {
1246 return Ok(Outcome::fail(
1247 FailureCode::Forbidden,
1248 format!("{}\n", refused.reason.unwrap_or_default()),
1249 ));
1250 }
1251 }
1252 let public = found.as_ref().is_some_and(|repo| !repo.is_private);
1253 let decision = g1t_contracts::scopes::decide_git(&access, write, public);
1254 if !decision.allowed {
1255 return Ok(Outcome::fail(
1256 FailureCode::Forbidden,
1257 format!("{}\n", decision.reason.unwrap_or_default()),
1258 ));
1259 }
1260 if !write && !access.allows(g1t_contracts::scopes::Scope::CodeRead) {
1261 a.viewer = None;
1262 }
1263 }
1264
1265 // Anonymous callers are asked to authenticate whether or not the repo
1266 // exists, so private repos cannot be told apart from missing ones.
1267 let denied = || match &a.viewer {
1268 Some(_) => not_found(),
1269 None => Outcome::fail(FailureCode::Unauthenticated, "Authentication required."),
1270 };
1271 // An agent's token works through the API only: its sandbox has its
1272 // own way to push, to its own pull request.
1273 if a.viewer.as_ref().is_some_and(|user| user.kind == PrincipalKind::Agent) {
1274 return Ok(Outcome::fail(
1275 FailureCode::Forbidden,
1276 "A g1t agent's token cannot be used with git.",
1277 ));
1278 }
1279 if let (true, Some(user)) = (write, &a.viewer)
1280 && !user.verified
1281 {
1282 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1283 }
1284 let repo = match found {
1285 Some(repo) => {
1286 let allowed = if write {
1287 can_write(&repo, &a.viewer)
1288 } else {
1289 self.may_read(&repo, &a.viewer).await?
1290 };
1291 if !allowed {
1292 return Ok(denied());
1293 }
1294 // An archived repository, or a pull request's copy of one,
1295 // is read-only.
1296 if write {
1297 let archived = match &repo.fork_of {
1298 Some(source) => self.registry.by_id(source).await?,
1299 None => Some(repo.clone()),
1300 };
1301 match archived {
1302 Some(source) => {
1303 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1304 return Ok(Outcome::fail(code, format!("{message}\n")));
1305 }
1306 }
1307 // The repository it was copied from is deleted.
1308 None => return Ok(denied()),
1309 }
1310 }
1311 repo
1312 }
1313 None => {
1314 // Push to create, in a workspace the pusher belongs to.
1315 let owner = a
1316 .viewer
1317 .as_ref()
1318 .filter(|user| write && user.is_member(&a.path.namespace.to_lowercase()));
1319 let Some(owner) = owner else {
1320 return Ok(denied());
1321 };
1322 let created = self.create(push_to_create(owner, &a.path)).await?;
1323 match created {
1324 Outcome::Ok(repo) => repo,
1325 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
1326 }
1327 }
1328 };
1329 // A push, or a credential to push with, waits while the repository
1330 // moves between namespaces (moves.rs), and goes to where it is now.
1331 if write {
1332 return Ok(match self.unpaused(repo).await? {
1333 Ok(repo) => Outcome::Ok(repo),
1334 Err((code, message)) => Outcome::fail(code, format!("{message}\n")),
1335 });
1336 }
1337 Ok(Outcome::Ok(repo))
1338 }
1339
1340 async fn land(&self, a: LandArgs) -> Result<Outcome<Landed>> {
1341 let actor: Viewer = Some(a.actor.clone());
1342 let Some(source) = self.registry.by_id(&a.source_id).await? else {
1343 return Ok(not_found());
1344 };
1345 // A fork lands on the repository it came from; a branch on its own.
1346 let target = match &source.fork_of {
1347 Some(id) => self.registry.by_id(id).await?,
1348 None => Some(source.clone()),
1349 };
1350 let Some(target) = target.filter(|repo| can_read(repo, &actor)) else {
1351 return Ok(not_found());
1352 };
1353 if !registry::can(&target, &actor, Capability::Merge) {
1354 return Ok(Outcome::fail(
1355 FailureCode::Forbidden,
1356 access::needs(Capability::Merge, &format!("{}/{}", target.namespace, target.name)),
1357 ));
1358 }
1359 if !a.actor.verified {
1360 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1361 }
1362 if let Some((code, message)) = lifecycle::archived_refusal(&target) {
1363 return Ok(Outcome::fail(code, message));
1364 }
1365 // Moving between namespaces: wait for it (moves.rs). Both are read
1366 // again once it is done, for their new keys.
1367 let (source, target) = match (self.unpaused(source).await?, self.unpaused(target).await?) {
1368 (Ok(source), Ok(target)) => (source, target),
1369 (Err((code, message)), _) | (_, Err((code, message))) => return Ok(Outcome::fail(code, message)),
1370 };
1371
1372 let branch = &a.target_branch.clone().unwrap_or_else(|| target.default_branch.clone());
1373 let from_fork = source.id != target.id;
1374 let source_branch = match a.branch {
1375 Some(name) if !from_fork && name == *branch => {
1376 return Ok(Outcome::fail(
1377 FailureCode::Invalid,
1378 format!("{branch} cannot be merged into itself."),
1379 ));
1380 }
1381 Some(name) => name,
1382 None if from_fork => branch.clone(),
1383 None => {
1384 return Ok(Outcome::fail(
1385 FailureCode::Invalid,
1386 "Say which branch to merge.",
1387 ));
1388 }
1389 };
1390
1391 self.live(&source).await?;
1392 let source_git = self.store.open(&store_key(&source)).await?;
1393 let target_git = self.store.open(&store_key(&target)).await?;
1394 let history = source_git.log(&source_branch, MAX_ANCESTRY).await?;
1395 let Some(new) = history.first().map(|commit| commit.hash.clone()) else {
1396 return Ok(Outcome::fail(
1397 FailureCode::Conflict,
1398 "This pull request has no commits to merge.",
1399 ));
1400 };
1401 let old = target_git
1402 .log(branch, 1)
1403 .await?
1404 .into_iter()
1405 .next()
1406 .map(|commit| commit.hash);
1407
1408 if old.as_deref() == Some(new.as_str()) {
1409 return Ok(Outcome::Ok(Landed {
1410 commit: new,
1411 previous: None,
1412 }));
1413 }
1414 // Moving the branch to a commit that does not descend from its
1415 // current head would discard whatever landed in between.
1416 if let Some(old) = &old
1417 && !descends_from(&source_git, &history, old).await?
1418 {
1419 let remedy = if from_fork {
1420 format!("Pull {branch} into the pull request's fork, push, and merge again.")
1421 } else {
1422 format!("Merge {branch} into {source_branch}, push, and merge again.")
1423 };
1424 return Ok(Outcome::fail(
1425 FailureCode::Conflict,
1426 format!("{branch} has moved since this pull request was opened. {remedy}"),
1427 ));
1428 }
1429
1430 // For a branch the objects are already in the target; sending them
1431 // again is harmless and keeps one way of moving a ref.
1432 let source_access = source_git.access(Scope::Read).await?;
1433 let target_access = target_git.access(Scope::Write).await?;
1434 let pushed =
1435 land::fast_forward(&source_access, &target_access, branch, old.as_deref(), &new)
1436 .await?;
1437 self.refs_moved(&target.id).await;
1438 if let Err(reason) = pushed {
1439 // Most often another pull request landed between the check and the push.
1440 return Ok(Outcome::fail(
1441 FailureCode::Conflict,
1442 format!("{branch} could not be updated: {reason}"),
1443 ));
1444 }
1445 self.publish_push(
1446 &target,
1447 &format!("refs/heads/{branch}"),
1448 old.as_deref(),
1449 &new,
1450 Some(&a.actor),
1451 )
1452 .await?;
1453 Ok(Outcome::Ok(Landed {
1454 commit: new,
1455 previous: old,
1456 }))
1457 }
1458
1459 async fn compare(&self, a: CompareArgs) -> Result<Outcome<Comparison>> {
1460 let Some(repo) = self
1461 .visible(self.registry.by_id(&a.repo_id).await?, &a.viewer)
1462 .await?
1463 else {
1464 return Ok(not_found());
1465 };
1466 let git = self.read_git(&repo).await?;
1467 let head_ref = a.head.as_deref().unwrap_or(&repo.default_branch);
1468 // A pull request into another branch is compared from where it
1469 // left that branch.
1470 let base_branch = a.base_branch.clone();
1471 // The head's history is only searched when the base is worked out
1472 // from another branch.
1473 let depth = if a.base.is_some() || is_commit_hash(head_ref) { 1 } else { MAX_ANCESTRY };
1474 let history = git.log(head_ref, depth).await?;
1475 let Some(head) = history.first() else {
1476 return Ok(Outcome::fail(
1477 FailureCode::Conflict,
1478 "There are no commits to compare.",
1479 ));
1480 };
1481
1482 // Where the head's history meets the default branch of `against`,
1483 // or the branch asked for.
1484 let shared_with = async |against: &Repo| -> Result<Option<String>> {
1485 let against_git = self.read_git(against).await?;
1486 let branch = base_branch.as_deref().unwrap_or(&against.default_branch);
1487 let shared: HashSet<String> = against_git
1488 .log(branch, MAX_ANCESTRY)
1489 .await?
1490 .into_iter()
1491 .map(|commit| commit.hash)
1492 .collect();
1493 nearest_ancestor_in(&git, &history, &shared).await
1494 };
1495 let base = match (a.base, &repo.fork_of) {
1496 (Some(base), _) => Some(base),
1497 // A fork is compared with the last commit it shares with the
1498 // repository it came from.
1499 (None, Some(target_id)) => match self.registry.by_id(target_id).await? {
1500 Some(target) => shared_with(&target).await?,
1501 None => None,
1502 },
1503 // A branch, with the point where it left the default branch.
1504 // A single commit, with its first parent.
1505 (None, None) if is_commit_hash(head_ref) => head.parents.first().cloned(),
1506 (None, None) if head_ref != base_branch.as_deref().unwrap_or(&repo.default_branch) => {
1507 shared_with(&repo).await?
1508 }
1509 (None, None) => head.parents.first().cloned(),
1510 };
1511 let base_tree = match &base {
1512 Some(base) => git
1513 .log(base, 1)
1514 .await?
1515 .into_iter()
1516 .next()
1517 .map(|commit| commit.tree_hash),
1518 None => None,
1519 };
1520 let (files, truncated) =
1521 diff::compare_trees(&git, base_tree.as_deref(), &head.tree_hash).await?;
1522 Ok(Outcome::Ok(Comparison {
1523 base,
1524 head: head.hash.clone(),
1525 files,
1526 truncated,
1527 }))
1528 }
1529
1530 /// Reports that `git_ref` of `repo` (a full ref) now points to `after`,
1531 /// moved by `actor` (marked when that was a workflow job's token).
1532 async fn publish_push(
1533 &self,
1534 repo: &Repo,
1535 git_ref: &str,
1536 before: Option<&str>,
1537 after: &str,
1538 actor: Option<&User>,
1539 ) -> Result<()> {
1540 let caused_by_job = actor.and_then(g1t_contracts::events::job_run_of).map(str::to_owned);
1541 self.publish_git_push(repo, git_ref, before, after, actor.map(|user| user.id.clone()), false, caused_by_job).await
1542 }
1543
1544 /// `publish_push`, saying whether the push reached the store without
1545 /// being scanned for secrets first.
1546 #[allow(clippy::too_many_arguments)]
1547 async fn publish_git_push(
1548 &self,
1549 repo: &Repo,
1550 git_ref: &str,
1551 before: Option<&str>,
1552 after: &str,
1553 actor: Option<String>,
1554 unscanned: bool,
1555 caused_by_job: Option<String>,
1556 ) -> Result<()> {
1557 self.publish(NewEvent {
1558 kind: "git.push",
1559 source: SOURCE,
1560 repo_id: Some(repo.id.clone()),
1561 actor,
1562 data: GitPush {
1563 repo_id: repo.id.clone(),
1564 git_ref: git_ref.to_owned(),
1565 before: before.map(str::to_owned),
1566 after: after.to_owned(),
1567 default_branch: git_ref.strip_prefix("refs/heads/")
1568 == Some(repo.default_branch.as_str()),
1569 unscanned,
1570 caused_by_job,
1571 },
1572 })
1573 .await
1574 }
1575
1576 /// Git over HTTPS. Only what decides the answer happens before it:
1577 /// the repository, who is asking and whether they may, the free
1578 /// workspace limits, push protection, and the store's own answer. The
1579 /// audit entry and what a push changed are recorded once git has its
1580 /// answer. Each answer says how long its steps took (`Server-Timing`).
1581 async fn git_http(&self, request: Request, env: &Env, ctx: &Context) -> Result<Response> {
1582 let mut timing = git_http::Timing::start();
1583 let Some(git) = git_http::parse(&request.url()?) else {
1584 return Response::error("Not found", 404);
1585 };
1586 let response = match self.answer_git(request, &git, env, ctx, &mut timing).await {
1587 Ok(response) => response,
1588 // The git store is busy: git hears when to try again.
1589 Err(error) => match resilience::busy(&error.to_string()) {
1590 Some(busy) => git_http::busy_response(busy)?,
1591 None => return Err(error),
1592 },
1593 };
1594 timing.apply(response)
1595 }
1596
1597 async fn answer_git(
1598 &self,
1599 request: Request,
1600 git: &git_http::GitRequest,
1601 env: &Env,
1602 ctx: &Context,
1603 timing: &mut git_http::Timing,
1604 ) -> Result<Response> {
1605 let write = git.service == GitService::ReceivePack;
1606 let get = request.method() == Method::Get;
1607 let identity = env.service("IDENTITY")?;
1608 // The repository and the caller's credentials, at once. A fetch may
1609 // go by the row as read a moment ago, for the same clone's next
1610 // request; a push always reads it. Anonymous callers cost nothing.
1611 let lookup = async {
1612 if write {
1613 self.registry.by_path(&git.path).await
1614 } else {
1615 self.registry.by_path_recent(&git.path).await
1616 }
1617 };
1618 let (found, viewer) =
1619 futures_util::future::join(lookup, git_http::viewer(&request, &identity)).await;
1620 let mut found = found?;
1621 timing.mark("repo");
1622 // A workspace alias staff set (identity's aliases.rs: `g1t` for
1623 // `flagon-io`) is answered in place, as the repository under the
1624 // workspace's slug: pushes and some clients do not follow
1625 // redirects. Everything after this sees only the workspace's slug.
1626 let aliased = match found {
1627 Some(_) => None,
1628 None => git_http::aliased(git, &identity).await?,
1629 };
1630 if let Some(aliased) = &aliased {
1631 found = if write {
1632 self.registry.by_path(&aliased.path).await?
1633 } else {
1634 self.registry.by_path_recent(&aliased.path).await?
1635 };
1636 timing.mark("alias");
1637 }
1638 let git = aliased.as_ref().unwrap_or(git);
1639 if found.is_none() {
1640 // A workspace that was renamed: git follows a redirect when it
1641 // first asks for refs, and uses the new address from then on.
1642 // A repository transferred to another workspace: the same, to
1643 // its new path. Fetches and pushes both follow either.
1644 let url = request.url()?;
1645 let (renamed, moved) = futures_util::future::join(
1646 git_http::renamed(&url, &identity),
1647 self.registry.resolve_moved(&git.path),
1648 )
1649 .await;
1650 timing.mark("moved");
1651 if let Some(location) = renamed? {
1652 return git_http::moved(&location, get);
1653 }
1654 if let Some(now) = moved?
1655 && let Some(location) = git_http::transferred(&url, &now)
1656 {
1657 return git_http::moved(&location, get);
1658 }
1659 }
1660 let viewer = viewer?;
1661 // A run credential is checked against its grants, then acts as the
1662 // person it works for. See run_access.rs.
1663 let (request, viewer, audit) = match self.admit_git(request, git, viewer, found.as_ref()).await? {
1664 run_access::Admitted::Go { request, viewer, entry } => (request, viewer, entry),
1665 run_access::Admitted::Refused(response) => return Ok(response),
1666 };
1667 let mut after = AfterGit {
1668 audit,
1669 status: 0,
1670 message: None,
1671 push: None,
1672 };
1673 let repo = match self.authorize_git(&git.path, &viewer, git.service, found).await? {
1674 Outcome::Ok(repo) => repo,
1675 refused => {
1676 let response = git_http::refuse(refused)?;
1677 after.ended(response.status_code(), None);
1678 after.spawn(env, ctx);
1679 return Ok(response);
1680 }
1681 };
1682 // A pull request's working copy removed after it closed is made
1683 // again before git uses it (forks.rs).
1684 self.live(&repo).await?;
1685 timing.mark("access");
1686 // Clones check out the default branch g1t keeps, which can have
1687 // changed since the store made the repository.
1688 let default_branch = repo.fork_of.is_none().then(|| repo.default_branch.clone());
1689 let key = store_key(&repo);
1690 let scope = if write { Scope::Write } else { Scope::Read };
1691 let mut request = request;
1692 let protocol = refs_cache::protocol(request.headers().get("git-protocol")?.as_deref());
1693 // A fetch's POST is read here, to tell an `ls-refs` from a fetch of
1694 // objects; the store would have it read in full anyway.
1695 let body = if !write && !get { Some(request.bytes().await?) } else { None };
1696 // What it asks the store, for the meters (meters.rs).
1697 let call = git_ops::classify(git.service, git.endpoint, get, body.as_deref());
1698 // Answers kept from the usual store may name refs the fallback
1699 // store does not have (fallback.rs): none are used, or kept.
1700 let fallback = self.store.on_fallback(&key);
1701 // An answer that lists refs may have been kept: see refs_cache.rs.
1702 let kept_key = refs_cache::kind(git, get, protocol, body.as_deref())
1703 .filter(|_| !fallback)
1704 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1705 .map(|(kind, version)| {
1706 refs_cache::Key::new(&repo.id, version, default_branch.as_deref(), protocol, &kind)
1707 });
1708 // A fresh clone's pack may have been kept too: see pack_cache.rs.
1709 // Under the same refs version, so never across a change to them.
1710 let pack_key = self
1711 .packs
1712 .as_ref()
1713 .filter(|_| !fallback)
1714 .and_then(|_| {
1715 let encoding = request.headers().get("content-encoding").ok().flatten();
1716 pack_cache::cacheable(git, get, protocol, encoding.as_deref(), body.as_deref())
1717 })
1718 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1719 .map(|(normalized, version)| pack_cache::Key::new(&repo.id, version, &normalized));
1720 // A kept answer and the free workspace limits, with a kept
1721 // credential looked up alongside. A kept answer goes back without
1722 // waiting for the credential, which it does not need.
1723 let ((answer, pack, limited), kept_access) = {
1724 let shared = self.shared.as_deref();
1725 let answer_and_limits = std::pin::pin!(futures_util::future::join3(
1726 async {
1727 match &kept_key {
1728 Some(kept_key) => refs_cache::get(shared, kept_key).await,
1729 None => None,
1730 }
1731 },
1732 async {
1733 match (&pack_key, self.packs.as_deref()) {
1734 (Some(pack_key), Some(packs)) => pack_cache::get(packs, pack_key).await,
1735 _ => None,
1736 }
1737 },
1738 self.git_limits(call, git, &repo, env),
1739 ));
1740 let kept_access = std::pin::pin!(self.store.kept_access(&key, scope));
1741 match futures_util::future::select(answer_and_limits, kept_access).await {
1742 futures_util::future::Either::Left((first, kept_access)) => {
1743 let answered = first.0.is_some() || first.1.is_some() || matches!(first.2, Ok(Some(_)) | Err(_));
1744 (first, if answered { None } else { kept_access.await })
1745 }
1746 futures_util::future::Either::Right((kept_access, first)) => (first.await, kept_access),
1747 }
1748 };
1749 timing.mark("kept");
1750 // A kept pack first: it never reaches the store, so it is never an
1751 // operation, and a free workspace past its operation cap still gets
1752 // it. (The other limits are a push's, and a pack is only a fetch.)
1753 if let Some(kept) = pack {
1754 timing.note("pack", "hit");
1755 let sent = body.as_ref().map_or(0, |body| body.len() as u64);
1756 meters::record(pack_cache::HIT, &key, sent, kept.size);
1757 after.ended(200, None);
1758 after.spawn(env, ctx);
1759 return kept.response();
1760 }
1761 if let Some((response, status, message)) = limited? {
1762 after.ended(status, Some(message.to_owned()));
1763 after.spawn(env, ctx);
1764 return Ok(response);
1765 }
1766 if pack_key.is_some() {
1767 timing.note("pack", "miss");
1768 }
1769 if let (Some((entry, found)), Some(kept_key)) = (answer, &kept_key) {
1770 timing.note("refs", found.as_str());
1771 if found == refs_cache::Found::Shared {
1772 let (kept_key, entry) = (kept_key.clone(), entry.clone());
1773 ctx.wait_until(async move { refs_cache::keep_in_colo(&kept_key, &entry).await });
1774 }
1775 // Never reached the store: never an operation.
1776 meters::record(call.cached_meter(), &key, 0, entry.body.len() as u64);
1777 after.ended(200, None);
1778 after.spawn(env, ctx);
1779 return entry.response();
1780 }
1781 if kept_key.is_some() {
1782 timing.note("refs", "miss");
1783 }
1784 // The store's credential: one made a moment ago, here or in another
1785 // isolate (see store.rs), or a new one.
1786 let access = match kept_access {
1787 Some((access, from)) => {
1788 timing.note("cred", from.as_str());
1789 access
1790 }
1791 None => {
1792 let access = self.store.mint_access(&key, scope).await?;
1793 timing.mark("mint");
1794 timing.note("cred", "mint");
1795 access
1796 }
1797 };
1798 // Should the store turn a kept credential down, a fetch's first
1799 // request is tried again with a new one; the requests after it then
1800 // have that one too.
1801 let again = if get { Some(request.clone()?) } else { None };
1802 // Push protection: a push that adds a secret is refused. See secret_scan.rs.
1803 let scan = async |body: &[u8]| self.protect(&repo, viewer.as_ref(), body).await;
1804 // Rulesets: what the rules of the branches and tags it changes
1805 // refuse is declined, saying which rule and why (rules.rs).
1806 let rules = async |head: &[u8], whole: bool| self.check_push(&repo, viewer.as_ref(), head, whole).await;
1807 // What a push may bring (pack_limits.rs): the repository's size is
1808 // its own and its pull requests' working copies'.
1809 let limits = if write && !get {
1810 git_http::PushLimits {
1811 held: self.held(&repo).await,
1812 repo_limit: self.repo_limit,
1813 large: self.large_pushes,
1814 ..git_http::PushLimits::default()
1815 }
1816 } else {
1817 git_http::PushLimits::default()
1818 };
1819 let mut outcome = git_http::forward(
1820 request,
1821 body,
1822 git,
1823 &access,
1824 rules,
1825 default_branch.as_deref(),
1826 limits,
1827 scan,
1828 )
1829 .await?;
1830 let turned_down = matches!(
1831 &outcome,
1832 git_http::Push::Forwarded(forwarded) if matches!(forwarded.response.status_code(), 401 | 403)
1833 );
1834 if turned_down {
1835 self.store.forget_access(&key).await;
1836 if let Some(again) = again {
1837 let access = self.store.mint_access(&key, scope).await?;
1838 let nothing = async |_: &[u8]| Ok(None);
1839 outcome = git_http::forward(
1840 again,
1841 None,
1842 git,
1843 &access,
1844 async |_: &[u8], _: bool| Ok(None),
1845 default_branch.as_deref(),
1846 git_http::PushLimits::default(),
1847 nothing,
1848 )
1849 .await?;
1850 }
1851 }
1852 let forwarded =
1853 match outcome {
1854 git_http::Push::Forwarded(forwarded) => forwarded,
1855 git_http::Push::Refused(response) => {
1856 after.ended(403, Some("The push was declined by rules.".to_owned()));
1857 after.spawn(env, ctx);
1858 return Ok(response);
1859 }
1860 git_http::Push::Blocked(response) => {
1861 after.ended(403, Some("The push adds a secret.".to_owned()));
1862 after.spawn(env, ctx);
1863 return Ok(response);
1864 }
1865 git_http::Push::Declined(response, reason) => {
1866 after.ended(403, Some(format!("The push was declined: {reason}.")));
1867 after.spawn(env, ctx);
1868 return Ok(response);
1869 }
1870 };
1871 if forwarded.from_store {
1872 let received = forwarded
1873 .response
1874 .headers()
1875 .get("content-length")?
1876 .and_then(|length| length.parse().ok())
1877 .unwrap_or(0);
1878 meters::record(call.meter(), &key, forwarded.sent, received);
1879 }
1880 timing.mark("store");
1881 let mut response = forwarded.response;
1882 let status = response.status_code();
1883 if write && !get {
1884 // A push: the store has moved its refs once it has answered in
1885 // full, so the answer is read before the change is recorded, and
1886 // only then goes back. Whoever fetches after it sees the push.
1887 let headers = response.headers().clone();
1888 headers.delete("content-length")?;
1889 let report = response.bytes().await?;
1890 self.refs_moved(&repo.id).await;
1891 timing.mark("refs");
1892 response = Response::from_bytes(report)?.with_headers(headers).with_status(status);
1893 } else if let (Some(kept_key), 200) = (&kept_key, status) {
1894 // A miss: this answer is kept for the next to ask.
1895 let headers = response.headers().clone();
1896 headers.delete("content-length")?;
1897 let body = response.bytes().await?;
1898 if let Some(content_type) = headers.get("content-type")? {
1899 let entry = refs_cache::Entry { content_type, body: body.clone() };
1900 if entry.keepable() {
1901 let shared = self.shared.clone();
1902 let kept_key = kept_key.clone();
1903 ctx.wait_until(async move { refs_cache::keep(shared.as_deref(), &kept_key, &entry).await });
1904 }
1905 }
1906 response = Response::from_bytes(body)?.with_headers(headers).with_status(status);
1907 } else if let (Some(pack_key), Some(packs), true) = (&pack_key, &self.packs, forwarded.from_store) {
1908 // A fresh clone the bucket did not have: counted, and its pack
1909 // kept as it streams to git, when it is a whole one.
1910 meters::record(pack_cache::MISS, &key, forwarded.sent, 0);
1911 if status == 200 {
1912 let store_key = key.clone();
1913 let measured = Box::new(move |bytes: u64| meters::record_bytes(pack_cache::MISS, &store_key, 0, bytes));
1914 let (teed, filling) = pack_cache::tee(response, packs.clone(), pack_key, measured)?;
1915 response = teed;
1916 if let Some(filling) = filling {
1917 let pack_key = pack_key.clone();
1918 ctx.wait_until(async move {
1919 let filled = filling.await;
1920 if !matches!(filled, pack_cache::Filled::Kept { .. } | pack_cache::Filled::Abandoned) {
1921 worker::console_warn!("pack {} not kept: {filled:?}", pack_key.as_str());
1922 }
1923 });
1924 }
1925 }
1926 }
1927 after.ended(status, None);
1928 if status == 200 && (forwarded.pack_bytes > 0 || !forwarded.pushed.is_empty()) {
1929 after.push = Some(PushDone {
1930 repo,
1931 pushed: forwarded.pushed,
1932 pack_bytes: forwarded.pack_bytes,
1933 caused_by_job: viewer.as_ref().and_then(g1t_contracts::events::job_run_of).map(str::to_owned),
1934 actor: viewer.map(|user: User| user.id),
1935 unscanned: forwarded.unscanned,
1936 });
1937 }
1938 after.spawn(env, ctx);
1939 Ok(response)
1940 }
1941
1942 /// The answer for a request a free workspace's limits stop, or a push
1943 /// to a full repository, with its status and reason for the audit log;
1944 /// `None` to go on.
1945 ///
1946 /// A clone, fetch or push is a git operation, which the git store
1947 /// charges g1t for: counted for billing once the answer has gone back
1948 /// (meters.rs), and a free workspace far past its share is slowed down
1949 /// rather than charged (see git_ops.rs). Whether it is past it is
1950 /// decided from counts this isolate already holds: the database is not
1951 /// asked on the way. A free workspace is never charged for private
1952 /// storage: once its private repositories hold the free amount, pushes
1953 /// to them stop, checked when a push begins so that git shows the
1954 /// reason. So do pushes to a repository at the store's size limit.
1955 async fn git_limits(
1956 &self,
1957 call: git_ops::GitCall,
1958 git: &git_http::GitRequest,
1959 repo: &Repo,
1960 env: &Env,
1961 ) -> Result<Option<(Response, u16, &'static str)>> {
1962 let namespace = git.path.namespace.to_lowercase();
1963 if meters::mapping_now().billable(call.meter()) > 0.0 {
1964 let now = now_ms();
1965 let hour = git_ops::hour_key(&rfc3339(now));
1966 let limits = git_ops::Limits::from_env(env);
1967 if let Some((month, hour_ops)) = git_ops::standing(&namespace, &hour, now)
1968 && git_ops::slow_down(month + 1, hour_ops + 1, limits.free_cap, limits.hourly)
1969 && git_ops::is_free_kept(env.service("BILLING").ok().as_ref(), &namespace).await
1970 {
1971 return Ok(Some((
1972 git_ops::too_many(&namespace, limits.free_cap, limits.hourly)?,
1973 429,
1974 "Too many git operations this hour.",
1975 )));
1976 }
1977 }
1978 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" {
1979 let held = self.held(repo).await;
1980 if held >= self.repo_limit {
1981 let message = format!(
1982 "{}/{} 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",
1983 repo.namespace,
1984 repo.name,
1985 pack_limits::megabytes(held)
1986 );
1987 return Ok(Some((Response::error(message, 403)?, 403, "The repository is full.")));
1988 }
1989 }
1990 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" && repo.is_private {
1991 let free = git_ops::free_private_bytes(env);
1992 let held = self.registry.private_bytes(&namespace).await.unwrap_or(0);
1993 if git_ops::storage_full(held, free)
1994 && git_ops::is_free(env.service("BILLING").ok().as_ref(), &namespace).await
1995 {
1996 return Ok(Some((
1997 git_ops::storage_full_response(&namespace, held, free)?,
1998 403,
1999 "Free private storage is full.",
2000 )));
2001 }
2002 }
2003 Ok(None)
2004 }
2005
2006 /// What a repository and its pull requests' working copies hold, as
2007 /// g1t counts it: read for a push's first request, kept a minute for
2008 /// the rest of it.
2009 async fn held(&self, repo: &Repo) -> u64 {
2010 let root = repo.fork_of.clone().unwrap_or_else(|| repo.id.clone());
2011 let now = now_ms();
2012 if let Some(held) = HELD.with(|held| held.borrow().get(&root, now)) {
2013 return held;
2014 }
2015 let held = self.registry.stored_bytes(&root).await.unwrap_or(0).max(0) as u64;
2016 HELD.with(|kept| kept.borrow_mut().put(root, held, now));
2017 held
2018 }
2019
2020 /// What a push changed, recorded once git has its answer.
2021 async fn record_push(&self, push: PushDone) -> Result<()> {
2022 let PushDone {
2023 repo,
2024 pushed,
2025 pack_bytes,
2026 actor,
2027 unscanned,
2028 caused_by_job,
2029 } = push;
2030 // What the push stored, for billing's storage meter. A failure only
2031 // leaves the count short.
2032 if pack_bytes > 0
2033 && let Err(error) = self.registry.add_stored_bytes(&repo, pack_bytes).await
2034 {
2035 worker::console_error!("stored bytes for {} not counted: {error}", repo.name);
2036 }
2037 if pushed.is_empty() {
2038 return Ok(());
2039 }
2040 // Artifacts' own push notifications are per repository, which does
2041 // not fit a repo per pull request, so the front end reports pushes
2042 // itself: one event for each branch that moved.
2043 let stored = self.store.open(&store_key(&repo)).await?;
2044 for pushed in &pushed {
2045 // The store can refuse one ref and accept another, so each
2046 // branch is checked against where it actually is. A tag the
2047 // store cannot read back is taken as pushed.
2048 let moved = match pushed.branch() {
2049 Some(branch) => stored
2050 .log(branch, 1)
2051 .await?
2052 .first()
2053 .is_some_and(|commit| commit.hash == pushed.after),
2054 None => stored.log(&pushed.git_ref, 1).await.map_or(true, |head| {
2055 head.first().is_none_or(|commit| commit.hash == pushed.after)
2056 }),
2057 };
2058 if moved {
2059 self.publish_git_push(
2060 &repo,
2061 &pushed.git_ref,
2062 pushed.before.as_deref(),
2063 &pushed.after,
2064 actor.clone(),
2065 unscanned,
2066 caused_by_job.clone(),
2067 )
2068 .await?;
2069 }
2070 }
2071 Ok(())
2072 }
2073}
2074
2075/// A push the store accepted, to be recorded once git has its answer.
2076struct PushDone {
2077 repo: Repo,
2078 pushed: Vec<git_http::Pushed>,
2079 pack_bytes: u64,
2080 actor: Option<String>,
2081 /// Too large to scan for secrets before it was stored.
2082 unscanned: bool,
2083 /// The run whose job's token pushed, if one did: its push starts no
2084 /// workflows.
2085 caused_by_job: Option<String>,
2086}
2087
2088/// What a git request leaves for after its answer: its audit entry, with
2089/// how the request ended, and what a push changed.
2090struct AfterGit {
2091 audit: Option<Box<g1t_contracts::audit::NewAuditEntry>>,
2092 status: u16,
2093 message: Option<String>,
2094 push: Option<PushDone>,
2095}
2096
2097impl AfterGit {
2098 fn ended(&mut self, status: u16, message: Option<String>) {
2099 self.status = status;
2100 self.message = message;
2101 }
2102
2103 /// Does the work once the response is on its way. A failure is logged:
2104 /// git has already been told how its request went.
2105 fn spawn(self, env: &Env, ctx: &Context) {
2106 if self.audit.is_none() && self.push.is_none() {
2107 return;
2108 }
2109 let env = env.clone();
2110 ctx.wait_until(async move {
2111 let repos = match service(&env) {
2112 Ok(repos) => repos,
2113 Err(error) => {
2114 worker::console_error!("git request not recorded: {error}");
2115 return;
2116 }
2117 };
2118 repos.finish_git(self.audit, self.status, self.message).await;
2119 if let Some(push) = self.push
2120 && let Err(error) = repos.record_push(push).await
2121 {
2122 worker::console_error!("push not recorded: {error}");
2123 }
2124 });
2125 }
2126}
2127
2128fn service(env: &Env) -> Result<Repos<ArtifactsStore>> {
2129 let shared = shared::Shared::from_env(env).map(Rc::new);
2130 Ok(Repos {
2131 registry: Registry { db: env.d1("DB")? },
2132 store: ArtifactsStore::new(env, shared.clone())?,
2133 shared,
2134 packs: pack_cache::Packs::from_env(env).map(Rc::new),
2135 events: env.service("EVENTS")?,
2136 security: env.service("SECURITY").ok(),
2137 billing: env.service("BILLING").ok(),
2138 identity: env.service("IDENTITY").ok(),
2139 work: env.service("WORK").ok(),
2140 free_private_bytes: git_ops::free_private_bytes(env),
2141 fork_days: forks::retention_days(env),
2142 repo_limit: env
2143 .var("REPO_STORAGE_LIMIT_BYTES")
2144 .ok()
2145 .and_then(|value| value.to_string().parse().ok())
2146 .unwrap_or(pack_limits::DEFAULT_REPO_LIMIT_BYTES),
2147 large_pushes: git_http::LargePushes::from_var(env.var("LARGE_PUSHES").ok().map(|value| value.to_string()).as_deref()),
2148 placement: shards::Placement::from_vars(
2149 env.var("ARTIFACTS_NEW_REPOS").ok().map(|value| value.to_string()).as_deref(),
2150 env.var("ARTIFACTS_EU_NAMESPACE").ok().map(|value| value.to_string()).as_deref(),
2151 ),
2152 limits: shards::limits(env.var("ARTIFACTS_NAMESPACE_LIMITS").ok().map(|value| value.to_string()).as_deref()),
2153 })
2154}
2155
2156/// Writes what this isolate metered once the answer has gone back, every
2157/// few seconds at most: now, or once it is due, waiting in this request's
2158/// `wait_until` so nothing counted is left for a request that may never
2159/// come (meters.rs).
2160fn flush_later(env: &Env, ctx: &Context) {
2161 let Some(wait) = meters::plan_flush() else {
2162 return;
2163 };
2164 if let Ok(db) = env.d1("DB") {
2165 ctx.wait_until(async move { meters::flush_after(&db, wait).await });
2166 }
2167}
2168
2169/// `/backups/<job id>/parts/<number>`: the job and the part's number.
2170fn backup_part_path(path: &str) -> Option<(String, u16)> {
2171 let rest = path.strip_prefix("/backups/")?;
2172 let (job, number) = rest.split_once("/parts/")?;
2173 let number = number.parse::<u16>().ok()?;
2174 (!job.is_empty() && !job.contains('/')).then(|| (job.to_owned(), number))
2175}
2176
2177fn backups_off<T>() -> Outcome<T> {
2178 Outcome::fail(FailureCode::Conflict, "Backups are off on this installation: it has no storage for them.")
2179}
2180
2181/// One part of a backup's bundle, with the job's token in its header.
2182async fn backup_part(request: &mut Request, env: &Env, repos: &Repos<ArtifactsStore>, job_id: String, number: u16) -> Result<Response> {
2183 let Some(blobs) = backups::storage(env) else {
2184 return reply(&backups_off::<()>());
2185 };
2186 let token = request.headers().get(g1t_contracts::backups::TOKEN_HEADER)?.unwrap_or_default();
2187 let bytes = request.bytes().await?;
2188 let job = g1t_contracts::backups::BackupJobArgs { job_id, token };
2189 reply(&backups::part(&repos.registry.db, &blobs, &job, number, bytes).await?)
2190}
2191
2192#[cfg(test)]
2193mod backup_path_tests {
2194 use super::backup_part_path;
2195
2196 #[test]
2197 fn a_part_is_named_by_its_job_and_number() {
2198 assert_eq!(backup_part_path("/backups/bkp_1/parts/3"), Some(("bkp_1".to_owned(), 3)));
2199 assert_eq!(backup_part_path("/backups/bkp_1/parts/x"), None);
2200 assert_eq!(backup_part_path("/backups//parts/1"), None);
2201 assert_eq!(backup_part_path("/acme/rocket.git/info/refs"), None);
2202 }
2203}
2204
2205/// Read methods whose answer is an `Outcome`: when the git store is busy,
2206/// the site is told so in words instead of failing the page.
2207const OUTCOME_READS: [&str; 6] = ["tree", "blob", "log", "branches", "blame", "compare"];
2208
2209#[event(fetch)]
2210async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
2211 let mut repos = service(&env)?;
2212 // A part of a backup's bundle, as the API passes it on from the
2213 // sandbox: bytes, not JSON (backups.rs).
2214 if request.method() == Method::Put
2215 && let Some((job_id, number)) = backup_part_path(&request.path())
2216 {
2217 let answered = backup_part(&mut request, &env, &repos, job_id, number).await;
2218 flush_later(&env, &ctx);
2219 return answered;
2220 }
2221 let Some(method) = rpc_method(&request) else {
2222 let answered = repos.git_http(request, &env, &ctx).await;
2223 flush_later(&env, &ctx);
2224 return answered;
2225 };
2226 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2227 // Git over HTTPS above always reads the primary.
2228 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2229 repos.registry.db = db;
2230 let body: serde_json::Value = request.json().await?;
2231
2232 let answered = async { match method.as_str() {
2233 "get" => reply(&repos.get(args(body)?).await?),
2234 "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
2235 "readable" => {
2236 let a: ReadableArgs = args(body)?;
2237 reply(&repos.registry.readable(&a.ids, &a.viewer).await?)
2238 }
2239 "public_namespaces" => {
2240 let a: PublicNamespacesArgs = args(body)?;
2241 reply(&repos.registry.public_namespaces(&a.owner_id).await?)
2242 }
2243 "path_by_id" => {
2244 let a: PathByIdArgs = args(body)?;
2245 reply(
2246 &repos
2247 .registry
2248 .by_id(&a.id)
2249 .await?
2250 .filter(|repo| repo.fork_of.is_none())
2251 .map(|repo| RepoPath {
2252 namespace: repo.namespace,
2253 name: repo.name,
2254 }),
2255 )
2256 }
2257 "list" => {
2258 let a: ListArgs = args(body)?;
2259 reply(
2260 &repos
2261 .registry
2262 .list(
2263 &a.viewer,
2264 a.query.as_deref(),
2265 a.namespace.as_deref(),
2266 a.member_only,
2267 )
2268 .await?,
2269 )
2270 }
2271 "create" => reply(&repos.create(args(body)?).await?),
2272 // Services only: a GitHub mirror catching up, or pushing out.
2273 "mirror" => reply(&repos.mirror(args(body)?).await?),
2274 "transfer" => reply(&repos.transfer(args(body)?).await?),
2275 // A repository's lifecycle: see lifecycle.rs.
2276 "delete" => reply(&repos.delete(args(body)?).await?),
2277 "deleted" => reply(&repos.deleted(args(body)?).await?),
2278 "restore" => reply(&repos.restore(args(body)?).await?),
2279 "purge" => reply(&repos.purge(args(body)?).await?),
2280 "purge_due" => reply(&repos.purge_due(args(body)?).await?),
2281 "rename" => reply(&repos.rename(args(body)?).await?),
2282 "archive" => reply(&repos.archive(args(body)?).await?),
2283 "set_visibility" => reply(&repos.set_visibility(args(body)?).await?),
2284 "set_default_branch" => reply(&repos.set_default_branch(args(body)?).await?),
2285 "rename_branch" => reply(&repos.rename_branch(args(body)?).await?),
2286 "resolve_branch" => reply(&repos.resolve_branch(args(body)?).await?),
2287 "status_by_id" => reply(&repos.status_by_id(args(body)?).await?),
2288 "resolve_path" => {
2289 let a: ResolvePathArgs = args(body)?;
2290 reply(&repos.registry.resolve_moved(&a.path).await?)
2291 }
2292 "namespace_count" => {
2293 let a: NamespaceCountArgs = args(body)?;
2294 reply(&repos.registry.count_in(&a.namespace).await?)
2295 }
2296 "update" => reply(&repos.update(args(body)?).await?),
2297 "tree" => reply(&repos.tree(args(body)?).await?),
2298 "blob" => reply(&repos.blob(args(body)?).await?),
2299 "log" => reply(&repos.log(args(body)?).await?),
2300 "blame" => reply(&repos.blame(args(body)?).await?),
2301 "fork_for_pull" => reply(&repos.fork_for_pull(args(body)?).await?),
2302 "git_access" => reply(&repos.git_access(args(body)?).await?),
2303 "branches" => reply(&repos.branches(args(body)?).await?),
2304 "last_commits" => reply(&repos.last_commits(args(body)?).await?),
2305 "tags" => reply(&repos.tags(args(body)?).await?),
2306 // The About: what the Files page shows beside the files (about.rs).
2307 // What is kept behind the head is worked out again after the answer.
2308 "about" => {
2309 let (answer, refresh) = repos.about(args(body)?).await?;
2310 about::refresh_later(&env, &ctx, refresh);
2311 reply(&answer)
2312 }
2313 "languages" => {
2314 let (answer, refresh) = repos.languages(args(body)?).await?;
2315 about::refresh_later(&env, &ctx, refresh);
2316 reply(&answer)
2317 }
2318 "contributors" => {
2319 let (answer, refresh) = repos.contributors(args(body)?).await?;
2320 about::refresh_later(&env, &ctx, refresh);
2321 reply(&answer)
2322 }
2323 "license" => {
2324 let (answer, refresh) = repos.license(args(body)?).await?;
2325 about::refresh_later(&env, &ctx, refresh);
2326 reply(&answer)
2327 }
2328 "stars" => reply(&repos.stars(args(body)?).await?),
2329 "star" => reply(&repos.star(args(body)?).await?),
2330 "stargazers" => reply(&repos.stargazers(args(body)?).await?),
2331 "starred" => reply(&repos.starred(args(body)?).await?),
2332 "releases" => reply(&repos.releases(args(body)?).await?),
2333 "release" => reply(&repos.release(args(body)?).await?),
2334 "create_release" => reply(&repos.create_release(args(body)?).await?),
2335 "update_release" => reply(&repos.update_release(args(body)?).await?),
2336 "delete_release" => reply(&repos.delete_release(args(body)?).await?),
2337 "head" => reply(&repos.head(args(body)?).await?),
2338 "behind" => reply(&repos.behind(args(body)?).await?),
2339 "divergence" => reply(&repos.divergence(args(body)?).await?),
2340 "land" => reply(&repos.land(args(body)?).await?),
2341 "update_pull_branch" => reply(&repos.update_pull_branch(args(body)?).await?),
2342 "delete_branch" => reply(&repos.delete_branch(args(body)?).await?),
2343 "commit_file" => reply(&repos.commit_file(args(body)?).await?),
2344 "compare" => reply(&repos.compare(args(body)?).await?),
2345 // Services only: a pull request's commits, as rules look at them (rules.rs).
2346 "inspect_commits" => reply(&repos.inspect_commits(args(body)?).await?),
2347 "scan_history" => reply(&repos.scan_history(args(body)?).await?),
2348 "find_lockfiles" => reply(&repos.find_lockfiles(args(body)?).await?),
2349 "match_pattern" => reply(&repos.match_pattern(args(body)?).await?),
2350 "check_secret" => reply(&repos.check_secret(args(body)?).await?),
2351 "list_files" => reply(&repos.list_files(args(body)?).await?),
2352 "changed_files" => reply(&repos.changed_files(args(body)?).await?),
2353 "read_blobs" => reply(&repos.read_blobs(args(body)?).await?),
2354 // Services only: what the Composer registry builds packages from.
2355 "refs" => reply(&repos.refs_of(args(body)?).await?),
2356 "raw_file" => reply(&repos.raw_file(args(body)?).await?),
2357 "raw_blobs" => reply(&repos.raw_blobs(args(body)?).await?),
2358 "visibility" => {
2359 let a: g1t_contracts::repos::VisibilityArgs = args(body)?;
2360 reply(&repos.registry.visibility(&a.paths).await?)
2361 }
2362 "storage" => reply(&repos.registry.storage().await?),
2363 "git_operations" => {
2364 let a: GitOperationsArgs = args(body)?;
2365 reply(&git_ops::totals(&repos.registry.db, &a.month, a.since.as_deref(), a.namespace.as_deref().map(str::to_lowercase).as_deref()).await?)
2366 }
2367 "all_ids" => {
2368 let a: AllIdsArgs = args(body)?;
2369 let limit = a.limit.clamp(1, 500);
2370 let ids = repos.registry.ids_after(a.after.as_deref(), limit).await?;
2371 let next = (ids.len() == limit as usize).then(|| ids.last().cloned()).flatten();
2372 reply(&IdPage { ids, next })
2373 }
2374 // The raw meters of the git store, for reconciling with Cloudflare
2375 // (meters.rs, scripts/ops/artifacts-usage.mjs).
2376 "artifacts_usage" => {
2377 let a: meters::UsageArgs = args(body)?;
2378 reply(&meters::usage(&repos.registry.db, &a).await?)
2379 }
2380 "operation_mapping" => reply(&meters::read_mapping(&repos.registry.db).await?),
2381 // Billing: the workspace each pull request's working copy is counted
2382 // for, so Cloudflare's own count of `pulls--<id>` shares out too.
2383 "pull_owners" => {
2384 #[derive(serde::Deserialize)]
2385 struct PullOwnersArgs {
2386 pulls: Vec<String>,
2387 }
2388 let a: PullOwnersArgs = args(body)?;
2389 let pulls: Vec<String> = a.pulls.into_iter().take(500).collect();
2390 reply(&serde_json::json!({ "owners": meters::pull_owners(&repos.registry.db, &pulls).await? }))
2391 }
2392 // Services only: which meters are operations, changed without a deploy.
2393 "set_operation_mapping" => {
2394 let row: meters::MappingRow = args(body)?;
2395 meters::set_mapping(&repos.registry.db, &row, &rfc3339(now_ms())).await?;
2396 reply(&meters::read_mapping(&repos.registry.db).await?)
2397 }
2398 // Backups (backups.rs): the runner's sweep claims queued ones, and
2399 // each sandbox, through the API, asks for its job and says how it went.
2400 "claim_backups" => {
2401 let a: g1t_contracts::backups::ClaimBackupsArgs = args(body)?;
2402 let blobs = backups::storage(&env);
2403 reply(&backups::claim(&repos.registry.db, blobs.as_ref(), &a, now_ms()).await?)
2404 }
2405 "backup_spec" => match backups::storage(&env) {
2406 Some(blobs) => {
2407 let a: g1t_contracts::backups::BackupJobArgs = args(body)?;
2408 let every = backups::Settings::from_env(&env).full_every;
2409 reply(&backups::spec(&repos.registry, &blobs, &repos.store, &a, every, now_ms()).await?)
2410 }
2411 None => reply(&backups_off::<bool>()),
2412 },
2413 "backup_complete" => match backups::storage(&env) {
2414 Some(blobs) => reply(&backups::complete(&repos.registry, &blobs, &args(body)?, now_ms()).await?),
2415 None => reply(&backups_off::<bool>()),
2416 },
2417 "backup_fail" => match backups::storage(&env) {
2418 Some(blobs) => reply(&backups::fail(&repos.registry.db, &blobs, &args(body)?).await?),
2419 None => reply(&backups_off::<bool>()),
2420 },
2421 // How the git store has been answering, for the status page.
2422 "store_health" => {
2423 let a: meters::HealthArgs = args(body)?;
2424 reply(&meters::health(&repos.registry.db, &a).await?)
2425 }
2426 // Where repositories may be kept, for a workspace's settings.
2427 "storage_options" => reply(&repos.storage_options()),
2428 // Services and operators only: how each namespace stands, and
2429 // moving a repository between them (namespaces.rs, moves.rs).
2430 "namespaces" => reply(&repos.standings().await?),
2431 "move_repository" => reply(&repos.move_repository(args(body)?).await?),
2432 "repository_moves" => {
2433 let a: moves::ListMovesArgs = args(body)?;
2434 reply(&repos.registry.moves(a.limit.unwrap_or(50)).await?)
2435 }
2436 _ => Response::error("Unknown method", 404),
2437 } }
2438 .await;
2439 // The git store is busy: said in words, with when to try again.
2440 let answered = match answered {
2441 Err(error) => match resilience::busy(&error.to_string()) {
2442 Some(busy) if OUTCOME_READS.contains(&method.as_str()) => {
2443 reply(&Outcome::<()>::fail(FailureCode::Conflict, busy.message().trim()))
2444 }
2445 Some(busy) => {
2446 let response = Response::error(busy.message(), 503)?;
2447 response.headers().set("retry-after", &busy.retry_after.to_string())?;
2448 Ok(response)
2449 }
2450 None => Err(error),
2451 },
2452 answered => answered,
2453 };
2454 flush_later(&env, &ctx);
2455 served.finish(answered)
2456}
2457
2458/// The nightly cron in wrangler.jsonc: tonight's backups are queued.
2459const BACKUP_CRON: &str = "53 2 * * *";
2460
2461/// The hourly sweep: deleted repositories whose time to be restored has
2462/// passed are purged. See lifecycle.rs. And, at [`BACKUP_CRON`], the
2463/// repositories whose refs moved are queued for a backup (backups.rs).
2464#[event(scheduled)]
2465async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
2466 let repos = match service(&env) {
2467 Ok(repos) => repos,
2468 Err(error) => {
2469 worker::console_error!("repos: the sweep could not start: {error}");
2470 return;
2471 }
2472 };
2473 if event.cron() == BACKUP_CRON {
2474 let Some(blobs) = backups::storage(&env) else { return };
2475 match backups::nightly(&repos.registry.db, &blobs, backups::Settings::from_env(&env), now_ms()).await {
2476 Ok(night) => worker::console_log!("repos: queued {} backups, removed {} of purged repositories", night.queued, night.pruned),
2477 Err(error) => worker::console_error!("repos: backups could not be queued: {error}"),
2478 }
2479 return;
2480 }
2481 match repos.purge_due(PurgeDueArgs::default()).await {
2482 Ok(0) => {}
2483 Ok(count) => worker::console_log!("repos: purged {count} deleted repositories"),
2484 Err(error) => worker::console_error!("repos: the purge sweep failed: {error}"),
2485 }
2486 // Pull requests' working copies whose time has come (forks.rs).
2487 match repos.retire_due().await {
2488 Ok(0) => {}
2489 Ok(count) => worker::console_log!("repos: removed {count} pull request working copies"),
2490 Err(error) => worker::console_error!("repos: the working copy sweep failed: {error}"),
2491 }
2492 // Repositories moving between namespaces, and old copies (moves.rs).
2493 match repos.run_moves().await {
2494 Ok(0) => {}
2495 Ok(count) => worker::console_log!("repos: moved {count} repositories between namespaces"),
2496 Err(error) => worker::console_error!("repos: the move sweep failed: {error}"),
2497 }
2498 meters::flush(&repos.registry.db).await;
2499}
2500
2501/// Events from the bus. A workspace's rename: its repositories move to the
2502/// workspace's current slug, asked of identity by id, so a repeated or late
2503/// delivery lands in the same place; their git store keys stay as they
2504/// were. A workspace's deletion: its repositories are deleted with it,
2505/// restored with it, or purged with it.
2506#[event(queue)]
2507async fn queue(batch: MessageBatch<Event>, env: Env, ctx: Context) -> Result<()> {
2508 let registry = Registry { db: env.d1("DB")? };
2509 let identity = env.service("IDENTITY")?;
2510 let handled = handle_events(&batch, &env, &registry, &identity).await;
2511 flush_later(&env, &ctx);
2512 handled
2513}
2514
2515async fn handle_events(batch: &MessageBatch<Event>, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
2516 for message in batch.messages()? {
2517 let event = message.body();
2518 // A pull request merged, closed or reopened: its working copy is
2519 // kept or let go (forks.rs).
2520 if let Some(change) = forks::pull_change(&event.kind) {
2521 let Some(pull_id) = forks::pull_id_of(&event.data) else {
2522 worker::console_error!("{} {} names no pull request", event.kind, event.id);
2523 continue;
2524 };
2525 let repos = service(env)?;
2526 match change {
2527 forks::PullChange::Settled => repos.pull_settled(&pull_id).await?,
2528 forks::PullChange::Reopened => repos.pull_reopened(&pull_id).await?,
2529 }
2530 continue;
2531 }
2532 // A workspace deleted, restored or purged: its repositories go with
2533 // it, come back with it, or are purged with it (lifecycle.rs).
2534 if event.kind == "workspace.deleting" {
2535 match serde_json::from_value::<WorkspaceDeleting>(event.data.clone()) {
2536 Ok(deleting) => service(env)?.delete_with_workspace(&deleting, &protected_workspaces(env)).await?,
2537 Err(_) => worker::console_error!("workspace.deleting {} could not be read", event.id),
2538 }
2539 continue;
2540 }
2541 if event.kind == "workspace.restored" {
2542 match serde_json::from_value::<WorkspaceRestored>(event.data.clone()) {
2543 Ok(restored) => service(env)?.restore_with_workspace(&restored).await?,
2544 Err(_) => worker::console_error!("workspace.restored {} could not be read", event.id),
2545 }
2546 continue;
2547 }
2548 if event.kind == "workspace.deleted" {
2549 match serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) {
2550 Ok(deleted) => service(env)?.purge_workspace(&deleted, &protected_workspaces(env)).await?,
2551 Err(_) => worker::console_error!("workspace.deleted {} could not be read", event.id),
2552 }
2553 continue;
2554 }
2555 if event.kind != "workspace.renamed" {
2556 continue;
2557 }
2558 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
2559 worker::console_error!("workspace.renamed {} could not be read", event.id);
2560 continue;
2561 };
2562 let names: HashMap<String, String> = g1t_kit::call(
2563 identity,
2564 "usernames",
2565 &g1t_contracts::identity::UsernamesArgs {
2566 ids: vec![renamed.workspace_id.clone()],
2567 },
2568 )
2569 .await?;
2570 let current = names
2571 .get(&renamed.workspace_id)
2572 .cloned()
2573 .unwrap_or_else(|| renamed.to.clone());
2574 let left = registry
2575 .rename_namespace(&renamed.stale_slugs(&current), &current)
2576 .await?;
2577 if left > 0 {
2578 worker::console_error!(
2579 "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
2580 renamed.from,
2581 renamed.to
2582 );
2583 }
2584 }
2585 Ok(())
2586}
2587
2588/// The workspaces whose repositories never go with a deletion, whatever is
2589/// published: `PROTECTED_WORKSPACES` if set here, and Flagon's always.
2590fn protected_workspaces(env: &Env) -> Vec<String> {
2591 let configured = env.var("PROTECTED_WORKSPACES").ok().map(|v| v.to_string());
2592 g1t_contracts::identity::protected_names(configured.as_deref())
2593}
2594
2595/// The repository a push to a path that does not exist yet creates: private,
2596/// so nothing pushed by mistake is published. An owner makes it public on
2597/// purpose (`POST /repos/{owner}/{repo}/visibility`).
2598fn push_to_create(owner: &User, path: &RepoPath) -> CreateArgs {
2599 CreateArgs {
2600 owner: owner.clone(),
2601 namespace: path.namespace.clone(),
2602 name: path.name.clone(),
2603 description: None,
2604 is_private: true,
2605 import_url: None,
2606 import_token: None,
2607 }
2608}
2609
2610#[cfg(test)]
2611mod push_to_create_tests {
2612 use super::*;
2613
2614 #[test]
2615 fn a_pushed_repository_starts_private() {
2616 let owner: User = serde_json::from_value(serde_json::json!({ "id": "usr_1", "username": "ada" })).unwrap();
2617 let args = push_to_create(&owner, &RepoPath { namespace: "acme".into(), name: "site".into() });
2618 assert!(args.is_private);
2619 assert_eq!((args.namespace.as_str(), args.name.as_str()), ("acme", "site"));
2620 }
2621}