Skip to content

g1t/services/repos/src/lib.rs

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