Skip to content
2,542 linesCodeBlameRaw
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 found = found?;
1596 timing.mark("repo");
1597 if found.is_none() {
1598 // A workspace that was renamed: git follows a redirect when it
1599 // first asks for refs, and uses the new address from then on.
1600 // A repository transferred to another workspace: the same, to
1601 // its new path. Fetches and pushes both follow either.
1602 let url = request.url()?;
1603 let (renamed, moved) = futures_util::future::join(
1604 git_http::renamed(&url, &identity),
1605 self.registry.resolve_moved(&git.path),
1606 )
1607 .await;
1608 timing.mark("moved");
1609 if let Some(location) = renamed? {
1610 return git_http::moved(&location, get);
1611 }
1612 if let Some(now) = moved?
1613 && let Some(location) = git_http::transferred(&url, &now)
1614 {
1615 return git_http::moved(&location, get);
1616 }
1617 }
1618 let viewer = viewer?;
1619 // A run credential is checked against its grants, then acts as the
1620 // person it works for. See run_access.rs.
1621 let (request, viewer, audit) = match self.admit_git(request, git, viewer, found.as_ref()).await? {
1622 run_access::Admitted::Go { request, viewer, entry } => (request, viewer, entry),
1623 run_access::Admitted::Refused(response) => return Ok(response),
1624 };
1625 let mut after = AfterGit {
1626 audit,
1627 status: 0,
1628 message: None,
1629 push: None,
1630 };
1631 let repo = match self.authorize_git(&git.path, &viewer, git.service, found).await? {
1632 Outcome::Ok(repo) => repo,
1633 refused => {
1634 let response = git_http::refuse(refused)?;
1635 after.ended(response.status_code(), None);
1636 after.spawn(env, ctx);
1637 return Ok(response);
1638 }
1639 };
1640 // A pull request's working copy removed after it closed is made
1641 // again before git uses it (forks.rs).
1642 self.live(&repo).await?;
1643 timing.mark("access");
1644 // Clones check out the default branch g1t keeps, which can have
1645 // changed since the store made the repository.
1646 let default_branch = repo.fork_of.is_none().then(|| repo.default_branch.clone());
1647 let key = store_key(&repo);
1648 let scope = if write { Scope::Write } else { Scope::Read };
1649 let mut request = request;
1650 let protocol = refs_cache::protocol(request.headers().get("git-protocol")?.as_deref());
1651 // A fetch's POST is read here, to tell an `ls-refs` from a fetch of
1652 // objects; the store would have it read in full anyway.
1653 let body = if !write && !get { Some(request.bytes().await?) } else { None };
1654 // What it asks the store, for the meters (meters.rs).
1655 let call = git_ops::classify(git.service, git.endpoint, get, body.as_deref());
1656 // Answers kept from the usual store may name refs the fallback
1657 // store does not have (fallback.rs): none are used, or kept.
1658 let fallback = self.store.on_fallback(&key);
1659 // An answer that lists refs may have been kept: see refs_cache.rs.
1660 let kept_key = refs_cache::kind(git, get, protocol, body.as_deref())
1661 .filter(|_| !fallback)
1662 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1663 .map(|(kind, version)| {
1664 refs_cache::Key::new(&repo.id, version, default_branch.as_deref(), protocol, &kind)
1665 });
1666 // A fresh clone's pack may have been kept too: see pack_cache.rs.
1667 // Under the same refs version, so never across a change to them.
1668 let pack_key = self
1669 .packs
1670 .as_ref()
1671 .filter(|_| !fallback)
1672 .and_then(|_| {
1673 let encoding = request.headers().get("content-encoding").ok().flatten();
1674 pack_cache::cacheable(git, get, protocol, encoding.as_deref(), body.as_deref())
1675 })
1676 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1677 .map(|(normalized, version)| pack_cache::Key::new(&repo.id, version, &normalized));
1678 // A kept answer and the free workspace limits, with a kept
1679 // credential looked up alongside. A kept answer goes back without
1680 // waiting for the credential, which it does not need.
1681 let ((answer, pack, limited), kept_access) = {
1682 let shared = self.shared.as_deref();
1683 let answer_and_limits = std::pin::pin!(futures_util::future::join3(
1684 async {
1685 match &kept_key {
1686 Some(kept_key) => refs_cache::get(shared, kept_key).await,
1687 None => None,
1688 }
1689 },
1690 async {
1691 match (&pack_key, self.packs.as_deref()) {
1692 (Some(pack_key), Some(packs)) => pack_cache::get(packs, pack_key).await,
1693 _ => None,
1694 }
1695 },
1696 self.git_limits(call, git, &repo, env),
1697 ));
1698 let kept_access = std::pin::pin!(self.store.kept_access(&key, scope));
1699 match futures_util::future::select(answer_and_limits, kept_access).await {
1700 futures_util::future::Either::Left((first, kept_access)) => {
1701 let answered = first.0.is_some() || first.1.is_some() || matches!(first.2, Ok(Some(_)) | Err(_));
1702 (first, if answered { None } else { kept_access.await })
1703 }
1704 futures_util::future::Either::Right((kept_access, first)) => (first.await, kept_access),
1705 }
1706 };
1707 timing.mark("kept");
1708 // A kept pack first: it never reaches the store, so it is never an
1709 // operation, and a free workspace past its operation cap still gets
1710 // it. (The other limits are a push's, and a pack is only a fetch.)
1711 if let Some(kept) = pack {
1712 timing.note("pack", "hit");
1713 let sent = body.as_ref().map_or(0, |body| body.len() as u64);
1714 meters::record(pack_cache::HIT, &key, sent, kept.size);
1715 after.ended(200, None);
1716 after.spawn(env, ctx);
1717 return kept.response();
1718 }
1719 if let Some((response, status, message)) = limited? {
1720 after.ended(status, Some(message.to_owned()));
1721 after.spawn(env, ctx);
1722 return Ok(response);
1723 }
1724 if pack_key.is_some() {
1725 timing.note("pack", "miss");
1726 }
1727 if let (Some((entry, found)), Some(kept_key)) = (answer, &kept_key) {
1728 timing.note("refs", found.as_str());
1729 if found == refs_cache::Found::Shared {
1730 let (kept_key, entry) = (kept_key.clone(), entry.clone());
1731 ctx.wait_until(async move { refs_cache::keep_in_colo(&kept_key, &entry).await });
1732 }
1733 // Never reached the store: never an operation.
1734 meters::record(call.cached_meter(), &key, 0, entry.body.len() as u64);
1735 after.ended(200, None);
1736 after.spawn(env, ctx);
1737 return entry.response();
1738 }
1739 if kept_key.is_some() {
1740 timing.note("refs", "miss");
1741 }
1742 // The store's credential: one made a moment ago, here or in another
1743 // isolate (see store.rs), or a new one.
1744 let access = match kept_access {
1745 Some((access, from)) => {
1746 timing.note("cred", from.as_str());
1747 access
1748 }
1749 None => {
1750 let access = self.store.mint_access(&key, scope).await?;
1751 timing.mark("mint");
1752 timing.note("cred", "mint");
1753 access
1754 }
1755 };
1756 // Should the store turn a kept credential down, a fetch's first
1757 // request is tried again with a new one; the requests after it then
1758 // have that one too.
1759 let again = if get { Some(request.clone()?) } else { None };
1760 // Push protection: a push that adds a secret is refused. See secret_scan.rs.
1761 let scan = async |body: &[u8]| self.protect(&repo, viewer.as_ref(), body).await;
1762 // Rulesets: what the rules of the branches and tags it changes
1763 // refuse is declined, saying which rule and why (rules.rs).
1764 let rules = async |head: &[u8], whole: bool| self.check_push(&repo, viewer.as_ref(), head, whole).await;
1765 // What a push may bring (pack_limits.rs): the repository's size is
1766 // its own and its pull requests' working copies'.
1767 let limits = if write && !get {
1768 git_http::PushLimits {
1769 held: self.held(&repo).await,
1770 repo_limit: self.repo_limit,
1771 large: self.large_pushes,
1772 ..git_http::PushLimits::default()
1773 }
1774 } else {
1775 git_http::PushLimits::default()
1776 };
1777 let mut outcome = git_http::forward(
1778 request,
1779 body,
1780 git,
1781 &access,
1782 rules,
1783 default_branch.as_deref(),
1784 limits,
1785 scan,
1786 )
1787 .await?;
1788 let turned_down = matches!(
1789 &outcome,
1790 git_http::Push::Forwarded(forwarded) if matches!(forwarded.response.status_code(), 401 | 403)
1791 );
1792 if turned_down {
1793 self.store.forget_access(&key).await;
1794 if let Some(again) = again {
1795 let access = self.store.mint_access(&key, scope).await?;
1796 let nothing = async |_: &[u8]| Ok(None);
1797 outcome = git_http::forward(
1798 again,
1799 None,
1800 git,
1801 &access,
1802 async |_: &[u8], _: bool| Ok(None),
1803 default_branch.as_deref(),
1804 git_http::PushLimits::default(),
1805 nothing,
1806 )
1807 .await?;
1808 }
1809 }
1810 let forwarded =
1811 match outcome {
1812 git_http::Push::Forwarded(forwarded) => forwarded,
1813 git_http::Push::Refused(response) => {
1814 after.ended(403, Some("The push was declined by rules.".to_owned()));
1815 after.spawn(env, ctx);
1816 return Ok(response);
1817 }
1818 git_http::Push::Blocked(response) => {
1819 after.ended(403, Some("The push adds a secret.".to_owned()));
1820 after.spawn(env, ctx);
1821 return Ok(response);
1822 }
1823 git_http::Push::Declined(response, reason) => {
1824 after.ended(403, Some(format!("The push was declined: {reason}.")));
1825 after.spawn(env, ctx);
1826 return Ok(response);
1827 }
1828 };
1829 if forwarded.from_store {
1830 let received = forwarded
1831 .response
1832 .headers()
1833 .get("content-length")?
1834 .and_then(|length| length.parse().ok())
1835 .unwrap_or(0);
1836 meters::record(call.meter(), &key, forwarded.sent, received);
1837 }
1838 timing.mark("store");
1839 let mut response = forwarded.response;
1840 let status = response.status_code();
1841 if write && !get {
1842 // A push: the store has moved its refs once it has answered in
1843 // full, so the answer is read before the change is recorded, and
1844 // only then goes back. Whoever fetches after it sees the push.
1845 let headers = response.headers().clone();
1846 headers.delete("content-length")?;
1847 let report = response.bytes().await?;
1848 self.refs_moved(&repo.id).await;
1849 timing.mark("refs");
1850 response = Response::from_bytes(report)?.with_headers(headers).with_status(status);
1851 } else if let (Some(kept_key), 200) = (&kept_key, status) {
1852 // A miss: this answer is kept for the next to ask.
1853 let headers = response.headers().clone();
1854 headers.delete("content-length")?;
1855 let body = response.bytes().await?;
1856 if let Some(content_type) = headers.get("content-type")? {
1857 let entry = refs_cache::Entry { content_type, body: body.clone() };
1858 if entry.keepable() {
1859 let shared = self.shared.clone();
1860 let kept_key = kept_key.clone();
1861 ctx.wait_until(async move { refs_cache::keep(shared.as_deref(), &kept_key, &entry).await });
1862 }
1863 }
1864 response = Response::from_bytes(body)?.with_headers(headers).with_status(status);
1865 } else if let (Some(pack_key), Some(packs), true) = (&pack_key, &self.packs, forwarded.from_store) {
1866 // A fresh clone the bucket did not have: counted, and its pack
1867 // kept as it streams to git, when it is a whole one.
1868 meters::record(pack_cache::MISS, &key, forwarded.sent, 0);
1869 if status == 200 {
1870 let store_key = key.clone();
1871 let measured = Box::new(move |bytes: u64| meters::record_bytes(pack_cache::MISS, &store_key, 0, bytes));
1872 let (teed, filling) = pack_cache::tee(response, packs.clone(), pack_key, measured)?;
1873 response = teed;
1874 if let Some(filling) = filling {
1875 let pack_key = pack_key.clone();
1876 ctx.wait_until(async move {
1877 let filled = filling.await;
1878 if !matches!(filled, pack_cache::Filled::Kept { .. } | pack_cache::Filled::Abandoned) {
1879 worker::console_warn!("pack {} not kept: {filled:?}", pack_key.as_str());
1880 }
1881 });
1882 }
1883 }
1884 }
1885 after.ended(status, None);
1886 if status == 200 && (forwarded.pack_bytes > 0 || !forwarded.pushed.is_empty()) {
1887 after.push = Some(PushDone {
1888 repo,
1889 pushed: forwarded.pushed,
1890 pack_bytes: forwarded.pack_bytes,
1891 actor: viewer.map(|user: User| user.id),
1892 unscanned: forwarded.unscanned,
1893 });
1894 }
1895 after.spawn(env, ctx);
1896 Ok(response)
1897 }
1898
1899 /// The answer for a request a free workspace's limits stop, or a push
1900 /// to a full repository, with its status and reason for the audit log;
1901 /// `None` to go on.
1902 ///
1903 /// A clone, fetch or push is a git operation, which the git store
1904 /// charges g1t for: counted for billing once the answer has gone back
1905 /// (meters.rs), and a free workspace far past its share is slowed down
1906 /// rather than charged (see git_ops.rs). Whether it is past it is
1907 /// decided from counts this isolate already holds: the database is not
1908 /// asked on the way. A free workspace is never charged for private
1909 /// storage: once its private repositories hold the free amount, pushes
1910 /// to them stop, checked when a push begins so that git shows the
1911 /// reason. So do pushes to a repository at the store's size limit.
1912 async fn git_limits(
1913 &self,
1914 call: git_ops::GitCall,
1915 git: &git_http::GitRequest,
1916 repo: &Repo,
1917 env: &Env,
1918 ) -> Result<Option<(Response, u16, &'static str)>> {
1919 let namespace = git.path.namespace.to_lowercase();
1920 if meters::mapping_now().billable(call.meter()) > 0.0 {
1921 let now = now_ms();
1922 let hour = git_ops::hour_key(&rfc3339(now));
1923 let limits = git_ops::Limits::from_env(env);
1924 if let Some((month, hour_ops)) = git_ops::standing(&namespace, &hour, now)
1925 && git_ops::slow_down(month + 1, hour_ops + 1, limits.free_cap, limits.hourly)
1926 && git_ops::is_free_kept(env.service("BILLING").ok().as_ref(), &namespace).await
1927 {
1928 return Ok(Some((
1929 git_ops::too_many(&namespace, limits.free_cap, limits.hourly)?,
1930 429,
1931 "Too many git operations this hour.",
1932 )));
1933 }
1934 }
1935 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" {
1936 let held = self.held(repo).await;
1937 if held >= self.repo_limit {
1938 let message = format!(
1939 "{}/{} 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",
1940 repo.namespace,
1941 repo.name,
1942 pack_limits::megabytes(held)
1943 );
1944 return Ok(Some((Response::error(message, 403)?, 403, "The repository is full.")));
1945 }
1946 }
1947 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" && repo.is_private {
1948 let free = git_ops::free_private_bytes(env);
1949 let held = self.registry.private_bytes(&namespace).await.unwrap_or(0);
1950 if git_ops::storage_full(held, free)
1951 && git_ops::is_free(env.service("BILLING").ok().as_ref(), &namespace).await
1952 {
1953 return Ok(Some((
1954 git_ops::storage_full_response(&namespace, held, free)?,
1955 403,
1956 "Free private storage is full.",
1957 )));
1958 }
1959 }
1960 Ok(None)
1961 }
1962
1963 /// What a repository and its pull requests' working copies hold, as
1964 /// g1t counts it: read for a push's first request, kept a minute for
1965 /// the rest of it.
1966 async fn held(&self, repo: &Repo) -> u64 {
1967 let root = repo.fork_of.clone().unwrap_or_else(|| repo.id.clone());
1968 let now = now_ms();
1969 if let Some(held) = HELD.with(|held| held.borrow().get(&root, now)) {
1970 return held;
1971 }
1972 let held = self.registry.stored_bytes(&root).await.unwrap_or(0).max(0) as u64;
1973 HELD.with(|kept| kept.borrow_mut().put(root, held, now));
1974 held
1975 }
1976
1977 /// What a push changed, recorded once git has its answer.
1978 async fn record_push(&self, push: PushDone) -> Result<()> {
1979 let PushDone {
1980 repo,
1981 pushed,
1982 pack_bytes,
1983 actor,
1984 unscanned,
1985 } = push;
1986 // What the push stored, for billing's storage meter. A failure only
1987 // leaves the count short.
1988 if pack_bytes > 0
1989 && let Err(error) = self.registry.add_stored_bytes(&repo, pack_bytes).await
1990 {
1991 worker::console_error!("stored bytes for {} not counted: {error}", repo.name);
1992 }
1993 if pushed.is_empty() {
1994 return Ok(());
1995 }
1996 // Artifacts' own push notifications are per repository, which does
1997 // not fit a repo per pull request, so the front end reports pushes
1998 // itself: one event for each branch that moved.
1999 let stored = self.store.open(&store_key(&repo)).await?;
2000 for pushed in &pushed {
2001 // The store can refuse one ref and accept another, so each
2002 // branch is checked against where it actually is. A tag the
2003 // store cannot read back is taken as pushed.
2004 let moved = match pushed.branch() {
2005 Some(branch) => stored
2006 .log(branch, 1)
2007 .await?
2008 .first()
2009 .is_some_and(|commit| commit.hash == pushed.after),
2010 None => stored.log(&pushed.git_ref, 1).await.map_or(true, |head| {
2011 head.first().is_none_or(|commit| commit.hash == pushed.after)
2012 }),
2013 };
2014 if moved {
2015 self.publish_git_push(
2016 &repo,
2017 &pushed.git_ref,
2018 pushed.before.as_deref(),
2019 &pushed.after,
2020 actor.clone(),
2021 unscanned,
2022 )
2023 .await?;
2024 }
2025 }
2026 Ok(())
2027 }
2028}
2029
2030/// A push the store accepted, to be recorded once git has its answer.
2031struct PushDone {
2032 repo: Repo,
2033 pushed: Vec<git_http::Pushed>,
2034 pack_bytes: u64,
2035 actor: Option<String>,
2036 /// Too large to scan for secrets before it was stored.
2037 unscanned: bool,
2038}
2039
2040/// What a git request leaves for after its answer: its audit entry, with
2041/// how the request ended, and what a push changed.
2042struct AfterGit {
2043 audit: Option<Box<g1t_contracts::audit::NewAuditEntry>>,
2044 status: u16,
2045 message: Option<String>,
2046 push: Option<PushDone>,
2047}
2048
2049impl AfterGit {
2050 fn ended(&mut self, status: u16, message: Option<String>) {
2051 self.status = status;
2052 self.message = message;
2053 }
2054
2055 /// Does the work once the response is on its way. A failure is logged:
2056 /// git has already been told how its request went.
2057 fn spawn(self, env: &Env, ctx: &Context) {
2058 if self.audit.is_none() && self.push.is_none() {
2059 return;
2060 }
2061 let env = env.clone();
2062 ctx.wait_until(async move {
2063 let repos = match service(&env) {
2064 Ok(repos) => repos,
2065 Err(error) => {
2066 worker::console_error!("git request not recorded: {error}");
2067 return;
2068 }
2069 };
2070 repos.finish_git(self.audit, self.status, self.message).await;
2071 if let Some(push) = self.push
2072 && let Err(error) = repos.record_push(push).await
2073 {
2074 worker::console_error!("push not recorded: {error}");
2075 }
2076 });
2077 }
2078}
2079
2080fn service(env: &Env) -> Result<Repos<ArtifactsStore>> {
2081 let shared = shared::Shared::from_env(env).map(Rc::new);
2082 Ok(Repos {
2083 registry: Registry { db: env.d1("DB")? },
2084 store: ArtifactsStore::new(env, shared.clone())?,
2085 shared,
2086 packs: pack_cache::Packs::from_env(env).map(Rc::new),
2087 events: env.service("EVENTS")?,
2088 security: env.service("SECURITY").ok(),
2089 billing: env.service("BILLING").ok(),
2090 identity: env.service("IDENTITY").ok(),
2091 work: env.service("WORK").ok(),
2092 free_private_bytes: git_ops::free_private_bytes(env),
2093 fork_days: forks::retention_days(env),
2094 repo_limit: env
2095 .var("REPO_STORAGE_LIMIT_BYTES")
2096 .ok()
2097 .and_then(|value| value.to_string().parse().ok())
2098 .unwrap_or(pack_limits::DEFAULT_REPO_LIMIT_BYTES),
2099 large_pushes: git_http::LargePushes::from_var(env.var("LARGE_PUSHES").ok().map(|value| value.to_string()).as_deref()),
2100 placement: shards::Placement::from_vars(
2101 env.var("ARTIFACTS_NEW_REPOS").ok().map(|value| value.to_string()).as_deref(),
2102 env.var("ARTIFACTS_EU_NAMESPACE").ok().map(|value| value.to_string()).as_deref(),
2103 ),
2104 limits: shards::limits(env.var("ARTIFACTS_NAMESPACE_LIMITS").ok().map(|value| value.to_string()).as_deref()),
2105 })
2106}
2107
2108/// Writes what this isolate metered once the answer has gone back, every
2109/// few seconds at most: now, or once it is due, waiting in this request's
2110/// `wait_until` so nothing counted is left for a request that may never
2111/// come (meters.rs).
2112fn flush_later(env: &Env, ctx: &Context) {
2113 let Some(wait) = meters::plan_flush() else {
2114 return;
2115 };
2116 if let Ok(db) = env.d1("DB") {
2117 ctx.wait_until(async move { meters::flush_after(&db, wait).await });
2118 }
2119}
2120
2121/// `/backups/<job id>/parts/<number>`: the job and the part's number.
2122fn backup_part_path(path: &str) -> Option<(String, u16)> {
2123 let rest = path.strip_prefix("/backups/")?;
2124 let (job, number) = rest.split_once("/parts/")?;
2125 let number = number.parse::<u16>().ok()?;
2126 (!job.is_empty() && !job.contains('/')).then(|| (job.to_owned(), number))
2127}
2128
2129fn backups_off<T>() -> Outcome<T> {
2130 Outcome::fail(FailureCode::Conflict, "Backups are off on this installation: it has no storage for them.")
2131}
2132
2133/// One part of a backup's bundle, with the job's token in its header.
2134async fn backup_part(request: &mut Request, env: &Env, repos: &Repos<ArtifactsStore>, job_id: String, number: u16) -> Result<Response> {
2135 let Some(blobs) = backups::storage(env) else {
2136 return reply(&backups_off::<()>());
2137 };
2138 let token = request.headers().get(g1t_contracts::backups::TOKEN_HEADER)?.unwrap_or_default();
2139 let bytes = request.bytes().await?;
2140 let job = g1t_contracts::backups::BackupJobArgs { job_id, token };
2141 reply(&backups::part(&repos.registry.db, &blobs, &job, number, bytes).await?)
2142}
2143
2144#[cfg(test)]
2145mod backup_path_tests {
2146 use super::backup_part_path;
2147
2148 #[test]
2149 fn a_part_is_named_by_its_job_and_number() {
2150 assert_eq!(backup_part_path("/backups/bkp_1/parts/3"), Some(("bkp_1".to_owned(), 3)));
2151 assert_eq!(backup_part_path("/backups/bkp_1/parts/x"), None);
2152 assert_eq!(backup_part_path("/backups//parts/1"), None);
2153 assert_eq!(backup_part_path("/acme/rocket.git/info/refs"), None);
2154 }
2155}
2156
2157/// Read methods whose answer is an `Outcome`: when the git store is busy,
2158/// the site is told so in words instead of failing the page.
2159const OUTCOME_READS: [&str; 6] = ["tree", "blob", "log", "branches", "blame", "compare"];
2160
2161#[event(fetch)]
2162async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
2163 let mut repos = service(&env)?;
2164 // A part of a backup's bundle, as the API passes it on from the
2165 // sandbox: bytes, not JSON (backups.rs).
2166 if request.method() == Method::Put
2167 && let Some((job_id, number)) = backup_part_path(&request.path())
2168 {
2169 let answered = backup_part(&mut request, &env, &repos, job_id, number).await;
2170 flush_later(&env, &ctx);
2171 return answered;
2172 }
2173 let Some(method) = rpc_method(&request) else {
2174 let answered = repos.git_http(request, &env, &ctx).await;
2175 flush_later(&env, &ctx);
2176 return answered;
2177 };
2178 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2179 // Git over HTTPS above always reads the primary.
2180 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2181 repos.registry.db = db;
2182 let body: serde_json::Value = request.json().await?;
2183
2184 let answered = async { match method.as_str() {
2185 "get" => reply(&repos.get(args(body)?).await?),
2186 "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
2187 "readable" => {
2188 let a: ReadableArgs = args(body)?;
2189 reply(&repos.registry.readable(&a.ids, &a.viewer).await?)
2190 }
2191 "public_namespaces" => {
2192 let a: PublicNamespacesArgs = args(body)?;
2193 reply(&repos.registry.public_namespaces(&a.owner_id).await?)
2194 }
2195 "path_by_id" => {
2196 let a: PathByIdArgs = args(body)?;
2197 reply(
2198 &repos
2199 .registry
2200 .by_id(&a.id)
2201 .await?
2202 .filter(|repo| repo.fork_of.is_none())
2203 .map(|repo| RepoPath {
2204 namespace: repo.namespace,
2205 name: repo.name,
2206 }),
2207 )
2208 }
2209 "list" => {
2210 let a: ListArgs = args(body)?;
2211 reply(
2212 &repos
2213 .registry
2214 .list(
2215 &a.viewer,
2216 a.query.as_deref(),
2217 a.namespace.as_deref(),
2218 a.member_only,
2219 )
2220 .await?,
2221 )
2222 }
2223 "create" => reply(&repos.create(args(body)?).await?),
2224 // Services only: a GitHub mirror catching up, or pushing out.
2225 "mirror" => reply(&repos.mirror(args(body)?).await?),
2226 "transfer" => reply(&repos.transfer(args(body)?).await?),
2227 // A repository's lifecycle: see lifecycle.rs.
2228 "delete" => reply(&repos.delete(args(body)?).await?),
2229 "deleted" => reply(&repos.deleted(args(body)?).await?),
2230 "restore" => reply(&repos.restore(args(body)?).await?),
2231 "purge" => reply(&repos.purge(args(body)?).await?),
2232 "purge_due" => reply(&repos.purge_due(args(body)?).await?),
2233 "rename" => reply(&repos.rename(args(body)?).await?),
2234 "archive" => reply(&repos.archive(args(body)?).await?),
2235 "set_visibility" => reply(&repos.set_visibility(args(body)?).await?),
2236 "set_default_branch" => reply(&repos.set_default_branch(args(body)?).await?),
2237 "rename_branch" => reply(&repos.rename_branch(args(body)?).await?),
2238 "resolve_branch" => reply(&repos.resolve_branch(args(body)?).await?),
2239 "status_by_id" => reply(&repos.status_by_id(args(body)?).await?),
2240 "resolve_path" => {
2241 let a: ResolvePathArgs = args(body)?;
2242 reply(&repos.registry.resolve_moved(&a.path).await?)
2243 }
2244 "namespace_count" => {
2245 let a: NamespaceCountArgs = args(body)?;
2246 reply(&repos.registry.count_in(&a.namespace).await?)
2247 }
2248 "update" => reply(&repos.update(args(body)?).await?),
2249 "tree" => reply(&repos.tree(args(body)?).await?),
2250 "blob" => reply(&repos.blob(args(body)?).await?),
2251 "log" => reply(&repos.log(args(body)?).await?),
2252 "blame" => reply(&repos.blame(args(body)?).await?),
2253 "fork_for_pull" => reply(&repos.fork_for_pull(args(body)?).await?),
2254 "git_access" => reply(&repos.git_access(args(body)?).await?),
2255 "branches" => reply(&repos.branches(args(body)?).await?),
2256 "last_commits" => reply(&repos.last_commits(args(body)?).await?),
2257 "tags" => reply(&repos.tags(args(body)?).await?),
2258 "head" => reply(&repos.head(args(body)?).await?),
2259 "behind" => reply(&repos.behind(args(body)?).await?),
2260 "divergence" => reply(&repos.divergence(args(body)?).await?),
2261 "land" => reply(&repos.land(args(body)?).await?),
2262 "update_pull_branch" => reply(&repos.update_pull_branch(args(body)?).await?),
2263 "delete_branch" => reply(&repos.delete_branch(args(body)?).await?),
2264 "commit_file" => reply(&repos.commit_file(args(body)?).await?),
2265 "compare" => reply(&repos.compare(args(body)?).await?),
2266 // Services only: a pull request's commits, as rules look at them (rules.rs).
2267 "inspect_commits" => reply(&repos.inspect_commits(args(body)?).await?),
2268 "scan_history" => reply(&repos.scan_history(args(body)?).await?),
2269 "find_lockfiles" => reply(&repos.find_lockfiles(args(body)?).await?),
2270 "match_pattern" => reply(&repos.match_pattern(args(body)?).await?),
2271 "check_secret" => reply(&repos.check_secret(args(body)?).await?),
2272 "list_files" => reply(&repos.list_files(args(body)?).await?),
2273 "changed_files" => reply(&repos.changed_files(args(body)?).await?),
2274 "read_blobs" => reply(&repos.read_blobs(args(body)?).await?),
2275 // Services only: what the Composer registry builds packages from.
2276 "refs" => reply(&repos.refs_of(args(body)?).await?),
2277 "raw_file" => reply(&repos.raw_file(args(body)?).await?),
2278 "raw_blobs" => reply(&repos.raw_blobs(args(body)?).await?),
2279 "visibility" => {
2280 let a: g1t_contracts::repos::VisibilityArgs = args(body)?;
2281 reply(&repos.registry.visibility(&a.paths).await?)
2282 }
2283 "storage" => reply(&repos.registry.storage().await?),
2284 "git_operations" => {
2285 let a: GitOperationsArgs = args(body)?;
2286 reply(&git_ops::totals(&repos.registry.db, &a.month, a.since.as_deref(), a.namespace.as_deref().map(str::to_lowercase).as_deref()).await?)
2287 }
2288 "all_ids" => {
2289 let a: AllIdsArgs = args(body)?;
2290 let limit = a.limit.clamp(1, 500);
2291 let ids = repos.registry.ids_after(a.after.as_deref(), limit).await?;
2292 let next = (ids.len() == limit as usize).then(|| ids.last().cloned()).flatten();
2293 reply(&IdPage { ids, next })
2294 }
2295 // The raw meters of the git store, for reconciling with Cloudflare
2296 // (meters.rs, scripts/ops/artifacts-usage.mjs).
2297 "artifacts_usage" => {
2298 let a: meters::UsageArgs = args(body)?;
2299 reply(&meters::usage(&repos.registry.db, &a).await?)
2300 }
2301 "operation_mapping" => reply(&meters::read_mapping(&repos.registry.db).await?),
2302 // Billing: the workspace each pull request's working copy is counted
2303 // for, so Cloudflare's own count of `pulls--<id>` shares out too.
2304 "pull_owners" => {
2305 #[derive(serde::Deserialize)]
2306 struct PullOwnersArgs {
2307 pulls: Vec<String>,
2308 }
2309 let a: PullOwnersArgs = args(body)?;
2310 let pulls: Vec<String> = a.pulls.into_iter().take(500).collect();
2311 reply(&serde_json::json!({ "owners": meters::pull_owners(&repos.registry.db, &pulls).await? }))
2312 }
2313 // Services only: which meters are operations, changed without a deploy.
2314 "set_operation_mapping" => {
2315 let row: meters::MappingRow = args(body)?;
2316 meters::set_mapping(&repos.registry.db, &row, &rfc3339(now_ms())).await?;
2317 reply(&meters::read_mapping(&repos.registry.db).await?)
2318 }
2319 // Backups (backups.rs): the runner's sweep claims queued ones, and
2320 // each sandbox, through the API, asks for its job and says how it went.
2321 "claim_backups" => {
2322 let a: g1t_contracts::backups::ClaimBackupsArgs = args(body)?;
2323 let blobs = backups::storage(&env);
2324 reply(&backups::claim(&repos.registry.db, blobs.as_ref(), &a, now_ms()).await?)
2325 }
2326 "backup_spec" => match backups::storage(&env) {
2327 Some(blobs) => {
2328 let a: g1t_contracts::backups::BackupJobArgs = args(body)?;
2329 let every = backups::Settings::from_env(&env).full_every;
2330 reply(&backups::spec(&repos.registry, &blobs, &repos.store, &a, every, now_ms()).await?)
2331 }
2332 None => reply(&backups_off::<bool>()),
2333 },
2334 "backup_complete" => match backups::storage(&env) {
2335 Some(blobs) => reply(&backups::complete(&repos.registry, &blobs, &args(body)?, now_ms()).await?),
2336 None => reply(&backups_off::<bool>()),
2337 },
2338 "backup_fail" => match backups::storage(&env) {
2339 Some(blobs) => reply(&backups::fail(&repos.registry.db, &blobs, &args(body)?).await?),
2340 None => reply(&backups_off::<bool>()),
2341 },
2342 // How the git store has been answering, for the status page.
2343 "store_health" => {
2344 let a: meters::HealthArgs = args(body)?;
2345 reply(&meters::health(&repos.registry.db, &a).await?)
2346 }
2347 // Where repositories may be kept, for a workspace's settings.
2348 "storage_options" => reply(&repos.storage_options()),
2349 // Services and operators only: how each namespace stands, and
2350 // moving a repository between them (namespaces.rs, moves.rs).
2351 "namespaces" => reply(&repos.standings().await?),
2352 "move_repository" => reply(&repos.move_repository(args(body)?).await?),
2353 "repository_moves" => {
2354 let a: moves::ListMovesArgs = args(body)?;
2355 reply(&repos.registry.moves(a.limit.unwrap_or(50)).await?)
2356 }
2357 _ => Response::error("Unknown method", 404),
2358 } }
2359 .await;
2360 // The git store is busy: said in words, with when to try again.
2361 let answered = match answered {
2362 Err(error) => match resilience::busy(&error.to_string()) {
2363 Some(busy) if OUTCOME_READS.contains(&method.as_str()) => {
2364 reply(&Outcome::<()>::fail(FailureCode::Conflict, busy.message().trim()))
2365 }
2366 Some(busy) => {
2367 let response = Response::error(busy.message(), 503)?;
2368 response.headers().set("retry-after", &busy.retry_after.to_string())?;
2369 Ok(response)
2370 }
2371 None => Err(error),
2372 },
2373 answered => answered,
2374 };
2375 flush_later(&env, &ctx);
2376 served.finish(answered)
2377}
2378
2379/// The nightly cron in wrangler.jsonc: tonight's backups are queued.
2380const BACKUP_CRON: &str = "53 2 * * *";
2381
2382/// The hourly sweep: deleted repositories whose time to be restored has
2383/// passed are purged. See lifecycle.rs. And, at [`BACKUP_CRON`], the
2384/// repositories whose refs moved are queued for a backup (backups.rs).
2385#[event(scheduled)]
2386async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
2387 let repos = match service(&env) {
2388 Ok(repos) => repos,
2389 Err(error) => {
2390 worker::console_error!("repos: the sweep could not start: {error}");
2391 return;
2392 }
2393 };
2394 if event.cron() == BACKUP_CRON {
2395 let Some(blobs) = backups::storage(&env) else { return };
2396 match backups::nightly(&repos.registry.db, &blobs, backups::Settings::from_env(&env), now_ms()).await {
2397 Ok(night) => worker::console_log!("repos: queued {} backups, removed {} of purged repositories", night.queued, night.pruned),
2398 Err(error) => worker::console_error!("repos: backups could not be queued: {error}"),
2399 }
2400 return;
2401 }
2402 match repos.purge_due(PurgeDueArgs::default()).await {
2403 Ok(0) => {}
2404 Ok(count) => worker::console_log!("repos: purged {count} deleted repositories"),
2405 Err(error) => worker::console_error!("repos: the purge sweep failed: {error}"),
2406 }
2407 // Pull requests' working copies whose time has come (forks.rs).
2408 match repos.retire_due().await {
2409 Ok(0) => {}
2410 Ok(count) => worker::console_log!("repos: removed {count} pull request working copies"),
2411 Err(error) => worker::console_error!("repos: the working copy sweep failed: {error}"),
2412 }
2413 // Repositories moving between namespaces, and old copies (moves.rs).
2414 match repos.run_moves().await {
2415 Ok(0) => {}
2416 Ok(count) => worker::console_log!("repos: moved {count} repositories between namespaces"),
2417 Err(error) => worker::console_error!("repos: the move sweep failed: {error}"),
2418 }
2419 meters::flush(&repos.registry.db).await;
2420}
2421
2422/// Events from the bus. A workspace's rename: its repositories move to the
2423/// workspace's current slug, asked of identity by id, so a repeated or late
2424/// delivery lands in the same place; their git store keys stay as they
2425/// were. A workspace's deletion: its repositories are deleted with it,
2426/// restored with it, or purged with it.
2427#[event(queue)]
2428async fn queue(batch: MessageBatch<Event>, env: Env, ctx: Context) -> Result<()> {
2429 let registry = Registry { db: env.d1("DB")? };
2430 let identity = env.service("IDENTITY")?;
2431 let handled = handle_events(&batch, &env, &registry, &identity).await;
2432 flush_later(&env, &ctx);
2433 handled
2434}
2435
2436async fn handle_events(batch: &MessageBatch<Event>, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
2437 for message in batch.messages()? {
2438 let event = message.body();
2439 // A pull request merged, closed or reopened: its working copy is
2440 // kept or let go (forks.rs).
2441 if let Some(change) = forks::pull_change(&event.kind) {
2442 let Some(pull_id) = forks::pull_id_of(&event.data) else {
2443 worker::console_error!("{} {} names no pull request", event.kind, event.id);
2444 continue;
2445 };
2446 let repos = service(env)?;
2447 match change {
2448 forks::PullChange::Settled => repos.pull_settled(&pull_id).await?,
2449 forks::PullChange::Reopened => repos.pull_reopened(&pull_id).await?,
2450 }
2451 continue;
2452 }
2453 // A workspace deleted, restored or purged: its repositories go with
2454 // it, come back with it, or are purged with it (lifecycle.rs).
2455 if event.kind == "workspace.deleting" {
2456 match serde_json::from_value::<WorkspaceDeleting>(event.data.clone()) {
2457 Ok(deleting) => service(env)?.delete_with_workspace(&deleting, &protected_workspaces(env)).await?,
2458 Err(_) => worker::console_error!("workspace.deleting {} could not be read", event.id),
2459 }
2460 continue;
2461 }
2462 if event.kind == "workspace.restored" {
2463 match serde_json::from_value::<WorkspaceRestored>(event.data.clone()) {
2464 Ok(restored) => service(env)?.restore_with_workspace(&restored).await?,
2465 Err(_) => worker::console_error!("workspace.restored {} could not be read", event.id),
2466 }
2467 continue;
2468 }
2469 if event.kind == "workspace.deleted" {
2470 match serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) {
2471 Ok(deleted) => service(env)?.purge_workspace(&deleted, &protected_workspaces(env)).await?,
2472 Err(_) => worker::console_error!("workspace.deleted {} could not be read", event.id),
2473 }
2474 continue;
2475 }
2476 if event.kind != "workspace.renamed" {
2477 continue;
2478 }
2479 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
2480 worker::console_error!("workspace.renamed {} could not be read", event.id);
2481 continue;
2482 };
2483 let names: HashMap<String, String> = g1t_kit::call(
2484 identity,
2485 "usernames",
2486 &g1t_contracts::identity::UsernamesArgs {
2487 ids: vec![renamed.workspace_id.clone()],
2488 },
2489 )
2490 .await?;
2491 let current = names
2492 .get(&renamed.workspace_id)
2493 .cloned()
2494 .unwrap_or_else(|| renamed.to.clone());
2495 let left = registry
2496 .rename_namespace(&renamed.stale_slugs(&current), &current)
2497 .await?;
2498 if left > 0 {
2499 worker::console_error!(
2500 "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
2501 renamed.from,
2502 renamed.to
2503 );
2504 }
2505 }
2506 Ok(())
2507}
2508
2509/// The workspaces whose repositories never go with a deletion, whatever is
2510/// published: `PROTECTED_WORKSPACES` if set here, and Flagon's always.
2511fn protected_workspaces(env: &Env) -> Vec<String> {
2512 let configured = env.var("PROTECTED_WORKSPACES").ok().map(|v| v.to_string());
2513 g1t_contracts::identity::protected_names(configured.as_deref())
2514}
2515
2516/// The repository a push to a path that does not exist yet creates: private,
2517/// so nothing pushed by mistake is published. An owner makes it public on
2518/// purpose (`POST /repos/{owner}/{repo}/visibility`).
2519fn push_to_create(owner: &User, path: &RepoPath) -> CreateArgs {
2520 CreateArgs {
2521 owner: owner.clone(),
2522 namespace: path.namespace.clone(),
2523 name: path.name.clone(),
2524 description: None,
2525 is_private: true,
2526 import_url: None,
2527 import_token: None,
2528 }
2529}
2530
2531#[cfg(test)]
2532mod push_to_create_tests {
2533 use super::*;
2534
2535 #[test]
2536 fn a_pushed_repository_starts_private() {
2537 let owner: User = serde_json::from_value(serde_json::json!({ "id": "usr_1", "username": "ada" })).unwrap();
2538 let args = push_to_create(&owner, &RepoPath { namespace: "acme".into(), name: "site".into() });
2539 assert!(args.is_private);
2540 assert_eq!((args.namespace.as_str(), args.name.as_str()), ("acme", "site"));
2541 }
2542}