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