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