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