Skip to content
2,985 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, MessageExt, 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 !deletable_branch(&a.branch, a.head.as_deref()) {
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 if a.branch == repo.default_branch {
1240 return Ok(Outcome::fail(FailureCode::Forbidden, "The default branch is never deleted."));
1241 }
1242 let repo = match self.unpaused(repo).await? {
1243 Ok(repo) => repo,
1244 Err((code, message)) => return Ok(Outcome::fail(code, message)),
1245 };
1246 self.live(&repo).await?;
1247 let git = self.store.open(&store_key(&repo)).await?;
1248 let Some(old) = git
1249 .branches()
1250 .await?
1251 .into_iter()
1252 .find(|branch| branch.name == a.branch)
1253 .map(|branch| branch.hash)
1254 else {
1255 return Ok(Outcome::Ok(false));
1256 };
1257 // Moved since the caller looked: someone else's commits are on it.
1258 if a.head.as_deref().is_some_and(|head| head != old) {
1259 return Ok(Outcome::fail(FailureCode::Conflict, format!("{} moved, so it was left alone.", a.branch)));
1260 }
1261 let access = git.access(Scope::Write).await?;
1262 let deleted = land::delete_ref(&access, &a.branch, &old).await?;
1263 self.refs_moved(&repo.id).await;
1264 if let Err(reason) = deleted {
1265 return Ok(Outcome::fail(
1266 FailureCode::Conflict,
1267 format!("{} could not be deleted: {reason}", a.branch),
1268 ));
1269 }
1270 Ok(Outcome::Ok(true))
1271 }
1272
1273 async fn fork_for_pull(&self, a: ForkArgs) -> Result<Outcome<Repo>> {
1274 let viewer = Some(a.actor.clone());
1275 let Some(source) = self
1276 .registry
1277 .by_id(&a.source_id)
1278 .await?
1279 .filter(|repo| can_read(repo, &viewer))
1280 else {
1281 return Ok(not_found());
1282 };
1283 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1284 return Ok(Outcome::fail(code, message));
1285 }
1286 // Its working copy is made in its namespace: not while it moves.
1287 let source = match self.unpaused(source).await? {
1288 Ok(source) => source,
1289 Err((code, message)) => return Ok(Outcome::fail(code, message)),
1290 };
1291 let now = now_ms();
1292 let fork = Repo {
1293 id: new_id("rep", now),
1294 namespace: PULLS_NAMESPACE.to_owned(),
1295 name: a.pull_id.clone(),
1296 description: None,
1297 // A fork is exactly as visible as the repo it came from.
1298 is_private: source.is_private,
1299 owner_id: a.actor.id.clone(),
1300 default_branch: source.default_branch.clone(),
1301 fork_of: Some(source.id.clone()),
1302 protected: false,
1303 created_at: rfc3339(now),
1304 topics: Vec::new(),
1305 website: None,
1306 archived_at: None,
1307 };
1308 // Artifacts forks within a namespace: the copy goes where its
1309 // repository is.
1310 let (namespace, _) = store::locate(&store_key(&source));
1311 self.registry
1312 .claim_store_key(&fork, Some(&namespace), &self.store.default_namespace())
1313 .await?;
1314 self.store
1315 .open(&store_key(&source))
1316 .await?
1317 .fork(&store_key(&fork))
1318 .await?;
1319 self.registry.insert(&fork).await?;
1320 self.publish(NewEvent {
1321 kind: "repo.forked",
1322 source: SOURCE,
1323 repo_id: Some(source.id.clone()),
1324 actor: Some(a.actor.id),
1325 data: RepoForked {
1326 repo_id: fork.id.clone(),
1327 source_repo_id: source.id,
1328 pull_id: a.pull_id,
1329 },
1330 })
1331 .await?;
1332 Ok(Outcome::Ok(fork))
1333 }
1334
1335 async fn git_access(&self, a: GitAccessArgs) -> Result<Outcome<GitAccess>> {
1336 let found = self.registry.by_path(&a.path).await?;
1337 Ok(match self.authorize_git(&a.path, &a.viewer, a.service, found).await? {
1338 Outcome::Ok(repo) => {
1339 self.live(&repo).await?;
1340 let write = a.service == GitService::ReceivePack;
1341 if write {
1342 // A push with this credential would not pass through
1343 // here, so nothing that lists the refs is kept until it
1344 // has expired (see refs_cache.rs).
1345 let until = now_ms() + store::CREDENTIAL_LIFE_MS + 60_000;
1346 if let Err(error) = self.registry.refs_open(&repo.id, until).await {
1347 // Before the column exists nothing is kept anyway.
1348 if registry::refs_state(&repo.id).is_some() {
1349 return Err(error);
1350 }
1351 }
1352 }
1353 let scope = if write { Scope::Write } else { Scope::Read };
1354 Outcome::Ok(self.store.handout(&store_key(&repo), scope).await?)
1355 }
1356 Outcome::Fail(failure) => Outcome::Fail(failure),
1357 })
1358 }
1359
1360 /// The repository at `path` (`found`, as just read), if the viewer may
1361 /// use `service` on it: fetch from it, or push to it. A push to a path
1362 /// with nothing there makes the repository, in a workspace the pusher
1363 /// belongs to.
1364 async fn authorize_git(
1365 &self,
1366 path: &RepoPath,
1367 viewer: &Viewer,
1368 service: GitService,
1369 found: Option<Repo>,
1370 ) -> Result<Outcome<Repo>> {
1371 let mut a = GitAccessArgs {
1372 path: path.clone(),
1373 viewer: viewer.clone(),
1374 service,
1375 };
1376 let write = a.service == GitService::ReceivePack;
1377 // An access token: pushing needs code:write, reading a private
1378 // repository code:read. A public repository reads as it would for
1379 // anyone. Which repositories a token reaches is its owner's, checked
1380 // below as for anyone.
1381 if let Some(access) = a.viewer.as_ref().and_then(|user| user.token.as_deref()).cloned() {
1382 // A workflow job's token, and a deploy key, reach their own
1383 // repository only; a job's also the working copies of that
1384 // repository's pull requests, where their heads are.
1385 let name = format!("{}/{}", path.namespace, path.name);
1386 let source = match found.as_ref().and_then(|repo| repo.fork_of.as_deref()) {
1387 Some(source_id) if g1t_contracts::scopes::decide_repo(&access, &name).is_some() => self
1388 .registry
1389 .by_id(source_id)
1390 .await?
1391 .map(|source| format!("{}/{}", source.namespace, source.name)),
1392 _ => None,
1393 };
1394 let public = found.as_ref().is_some_and(|repo| !repo.is_private);
1395 if let Some(why) = git_token_refusal(&access, &name, source.as_deref(), write, public, found.is_some()) {
1396 return Ok(Outcome::fail(FailureCode::Forbidden, format!("{why}\n")));
1397 }
1398 if !write && !access.allows(g1t_contracts::scopes::Scope::CodeRead) {
1399 a.viewer = None;
1400 }
1401 }
1402
1403 // Anonymous callers are asked to authenticate whether or not the repo
1404 // exists, so private repos cannot be told apart from missing ones.
1405 let denied = || match &a.viewer {
1406 Some(_) => not_found(),
1407 None => Outcome::fail(FailureCode::Unauthenticated, "Authentication required."),
1408 };
1409 // An agent's token works through the API only: its sandbox has its
1410 // own way to push, to its own pull request.
1411 if a.viewer.as_ref().is_some_and(|user| user.kind == PrincipalKind::Agent) {
1412 return Ok(Outcome::fail(
1413 FailureCode::Forbidden,
1414 "A g1t agent's token cannot be used with git.",
1415 ));
1416 }
1417 if let (true, Some(user)) = (write, &a.viewer)
1418 && !user.verified
1419 {
1420 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1421 }
1422 let repo = match found {
1423 Some(repo) => {
1424 let allowed = if write {
1425 can_write(&repo, &a.viewer)
1426 } else {
1427 self.may_read(&repo, &a.viewer).await?
1428 };
1429 if !allowed {
1430 return Ok(denied());
1431 }
1432 // An archived repository, or a pull request's copy of one,
1433 // is read-only.
1434 if write {
1435 let archived = match &repo.fork_of {
1436 Some(source) => self.registry.by_id(source).await?,
1437 None => Some(repo.clone()),
1438 };
1439 match archived {
1440 Some(source) => {
1441 if let Some((code, message)) = lifecycle::archived_refusal(&source) {
1442 return Ok(Outcome::fail(code, format!("{message}\n")));
1443 }
1444 }
1445 // The repository it was copied from is deleted.
1446 None => return Ok(denied()),
1447 }
1448 }
1449 repo
1450 }
1451 None => {
1452 // Push to create, in a workspace the pusher belongs to.
1453 let owner = a
1454 .viewer
1455 .as_ref()
1456 .filter(|user| write && user.is_member(&a.path.namespace.to_lowercase()));
1457 let Some(owner) = owner else {
1458 return Ok(denied());
1459 };
1460 let created = self.create(push_to_create(owner, &a.path)).await?;
1461 match created {
1462 Outcome::Ok(repo) => repo,
1463 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
1464 }
1465 }
1466 };
1467 // A push, or a credential to push with, waits while the repository
1468 // moves between namespaces (moves.rs), and goes to where it is now.
1469 if write {
1470 return Ok(match self.unpaused(repo).await? {
1471 Ok(repo) => Outcome::Ok(repo),
1472 Err((code, message)) => Outcome::fail(code, format!("{message}\n")),
1473 });
1474 }
1475 Ok(Outcome::Ok(repo))
1476 }
1477
1478 async fn land(&self, a: LandArgs) -> Result<Outcome<Landed>> {
1479 let actor: Viewer = Some(a.actor.clone());
1480 let Some(source) = self.registry.by_id(&a.source_id).await? else {
1481 return Ok(not_found());
1482 };
1483 // A fork lands on the repository it came from; a branch on its own.
1484 let target = match &source.fork_of {
1485 Some(id) => self.registry.by_id(id).await?,
1486 None => Some(source.clone()),
1487 };
1488 let Some(target) = target.filter(|repo| can_read(repo, &actor)) else {
1489 return Ok(not_found());
1490 };
1491 if !registry::can(&target, &actor, Capability::Merge) {
1492 return Ok(Outcome::fail(
1493 FailureCode::Forbidden,
1494 access::needs(Capability::Merge, &format!("{}/{}", target.namespace, target.name)),
1495 ));
1496 }
1497 if !a.actor.verified {
1498 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1499 }
1500 if let Some((code, message)) = lifecycle::archived_refusal(&target) {
1501 return Ok(Outcome::fail(code, message));
1502 }
1503 // Moving between namespaces: wait for it (moves.rs). Both are read
1504 // again once it is done, for their new keys.
1505 let (source, target) = match (self.unpaused(source).await?, self.unpaused(target).await?) {
1506 (Ok(source), Ok(target)) => (source, target),
1507 (Err((code, message)), _) | (_, Err((code, message))) => return Ok(Outcome::fail(code, message)),
1508 };
1509
1510 let branch = &a.target_branch.clone().unwrap_or_else(|| target.default_branch.clone());
1511 let from_fork = source.id != target.id;
1512 let source_branch = match a.branch {
1513 Some(name) if !from_fork && name == *branch => {
1514 return Ok(Outcome::fail(
1515 FailureCode::Invalid,
1516 format!("{branch} cannot be merged into itself."),
1517 ));
1518 }
1519 Some(name) => name,
1520 None if from_fork => branch.clone(),
1521 None => {
1522 return Ok(Outcome::fail(
1523 FailureCode::Invalid,
1524 "Say which branch to merge.",
1525 ));
1526 }
1527 };
1528
1529 self.live(&source).await?;
1530 let source_git = self.store.open(&store_key(&source)).await?;
1531 let target_git = self.store.open(&store_key(&target)).await?;
1532 let history = source_git.log(&source_branch, MAX_ANCESTRY).await?;
1533 let Some(new) = history.first().map(|commit| commit.hash.clone()) else {
1534 return Ok(Outcome::fail(
1535 FailureCode::Conflict,
1536 "This pull request has no commits to merge.",
1537 ));
1538 };
1539 let old = target_git
1540 .log(branch, 1)
1541 .await?
1542 .into_iter()
1543 .next()
1544 .map(|commit| commit.hash);
1545
1546 if old.as_deref() == Some(new.as_str()) {
1547 return Ok(Outcome::Ok(Landed {
1548 commit: new,
1549 previous: None,
1550 }));
1551 }
1552 // Moving the branch to a commit that does not descend from its
1553 // current head would discard whatever landed in between.
1554 if let Some(old) = &old
1555 && !descends_from(&source_git, &history, old).await?
1556 {
1557 let remedy = if from_fork {
1558 format!("Pull {branch} into the pull request's fork, push, and merge again.")
1559 } else {
1560 format!("Merge {branch} into {source_branch}, push, and merge again.")
1561 };
1562 return Ok(Outcome::fail(
1563 FailureCode::Conflict,
1564 format!("{branch} has moved since this pull request was opened. {remedy}"),
1565 ));
1566 }
1567
1568 // For a branch the objects are already in the target; sending them
1569 // again is harmless and keeps one way of moving a ref.
1570 let source_access = source_git.access(Scope::Read).await?;
1571 let target_access = target_git.access(Scope::Write).await?;
1572 let pushed =
1573 land::fast_forward(&source_access, &target_access, branch, old.as_deref(), &new)
1574 .await?;
1575 self.refs_moved(&target.id).await;
1576 if let Err(reason) = pushed {
1577 // Most often another pull request landed between the check and the push.
1578 return Ok(Outcome::fail(
1579 FailureCode::Conflict,
1580 format!("{branch} could not be updated: {reason}"),
1581 ));
1582 }
1583 self.publish_push(
1584 &target,
1585 &format!("refs/heads/{branch}"),
1586 old.as_deref(),
1587 &new,
1588 Some(&a.actor),
1589 )
1590 .await?;
1591 Ok(Outcome::Ok(Landed {
1592 commit: new,
1593 previous: old,
1594 }))
1595 }
1596
1597 async fn compare(&self, a: CompareArgs) -> Result<Outcome<Comparison>> {
1598 let Some(repo) = self
1599 .visible(self.registry.by_id(&a.repo_id).await?, &a.viewer)
1600 .await?
1601 else {
1602 return Ok(not_found());
1603 };
1604 let git = self.read_git(&repo).await?;
1605 let head_ref = a.head.as_deref().unwrap_or(&repo.default_branch);
1606 // A pull request into another branch is compared from where it
1607 // left that branch.
1608 let base_branch = a.base_branch.clone();
1609 // The head's history is only searched when the base is worked out
1610 // from another branch.
1611 let depth = if a.base.is_some() || is_commit_hash(head_ref) { 1 } else { MAX_ANCESTRY };
1612 let history = git.log(head_ref, depth).await?;
1613 let Some(head) = history.first() else {
1614 return Ok(Outcome::fail(
1615 FailureCode::Conflict,
1616 "There are no commits to compare.",
1617 ));
1618 };
1619
1620 // Where the head's history meets the default branch of `against`,
1621 // or the branch asked for.
1622 let shared_with = async |against: &Repo| -> Result<Option<String>> {
1623 let against_git = self.read_git(against).await?;
1624 let branch = base_branch.as_deref().unwrap_or(&against.default_branch);
1625 let shared: HashSet<String> = against_git
1626 .log(branch, MAX_ANCESTRY)
1627 .await?
1628 .into_iter()
1629 .map(|commit| commit.hash)
1630 .collect();
1631 nearest_ancestor_in(&git, &history, &shared).await
1632 };
1633 let base = match (a.base, &repo.fork_of) {
1634 (Some(base), _) => Some(base),
1635 // A fork is compared with the last commit it shares with the
1636 // repository it came from.
1637 (None, Some(target_id)) => match self.registry.by_id(target_id).await? {
1638 Some(target) => shared_with(&target).await?,
1639 None => None,
1640 },
1641 // A branch, with the point where it left the default branch.
1642 // A single commit, with its first parent.
1643 (None, None) if is_commit_hash(head_ref) => head.parents.first().cloned(),
1644 (None, None) if head_ref != base_branch.as_deref().unwrap_or(&repo.default_branch) => {
1645 shared_with(&repo).await?
1646 }
1647 (None, None) => head.parents.first().cloned(),
1648 };
1649 let base_tree = match &base {
1650 Some(base) => git
1651 .log(base, 1)
1652 .await?
1653 .into_iter()
1654 .next()
1655 .map(|commit| commit.tree_hash),
1656 None => None,
1657 };
1658 let (files, truncated) =
1659 diff::compare_trees(&git, base_tree.as_deref(), &head.tree_hash).await?;
1660 Ok(Outcome::Ok(Comparison {
1661 base,
1662 head: head.hash.clone(),
1663 files,
1664 truncated,
1665 }))
1666 }
1667
1668 /// Reports that `git_ref` of `repo` (a full ref) now points to `after`,
1669 /// moved by `actor` (marked when that was a workflow job's token).
1670 async fn publish_push(
1671 &self,
1672 repo: &Repo,
1673 git_ref: &str,
1674 before: Option<&str>,
1675 after: &str,
1676 actor: Option<&User>,
1677 ) -> Result<()> {
1678 let caused_by_job = actor.and_then(g1t_contracts::events::job_run_of).map(str::to_owned);
1679 self.publish_git_push(repo, git_ref, before, after, actor.map(|user| user.id.clone()), false, caused_by_job).await
1680 }
1681
1682 /// `publish_push`, saying whether the push reached the store without
1683 /// being scanned for secrets first.
1684 #[allow(clippy::too_many_arguments)]
1685 async fn publish_git_push(
1686 &self,
1687 repo: &Repo,
1688 git_ref: &str,
1689 before: Option<&str>,
1690 after: &str,
1691 actor: Option<String>,
1692 unscanned: bool,
1693 caused_by_job: Option<String>,
1694 ) -> Result<()> {
1695 self.publish(NewEvent {
1696 kind: "git.push",
1697 source: SOURCE,
1698 repo_id: Some(repo.id.clone()),
1699 actor,
1700 data: GitPush {
1701 repo_id: repo.id.clone(),
1702 git_ref: git_ref.to_owned(),
1703 before: before.map(str::to_owned),
1704 after: after.to_owned(),
1705 default_branch: git_ref.strip_prefix("refs/heads/")
1706 == Some(repo.default_branch.as_str()),
1707 unscanned,
1708 caused_by_job,
1709 },
1710 })
1711 .await
1712 }
1713
1714 /// Git over HTTPS. Only what decides the answer happens before it:
1715 /// the repository, who is asking and whether they may, the free
1716 /// workspace limits, push protection, and the store's own answer. The
1717 /// audit entry and what a push changed are recorded once git has its
1718 /// answer. Each answer says how long its steps took (`Server-Timing`).
1719 async fn git_http(&self, request: Request, env: &Env, ctx: &Context) -> Result<Response> {
1720 let mut timing = git_http::Timing::start();
1721 let Some(git) = git_http::parse(&request.url()?) else {
1722 return Response::error("Not found", 404);
1723 };
1724 let response = match self.answer_git(request, &git, env, ctx, &mut timing).await {
1725 Ok(response) => response,
1726 // The git store is busy: git hears when to try again.
1727 Err(error) => match resilience::busy(&error.to_string()) {
1728 Some(busy) => git_http::busy_response(busy)?,
1729 None => return Err(error),
1730 },
1731 };
1732 timing.apply(response)
1733 }
1734
1735 async fn answer_git(
1736 &self,
1737 request: Request,
1738 git: &git_http::GitRequest,
1739 env: &Env,
1740 ctx: &Context,
1741 timing: &mut git_http::Timing,
1742 ) -> Result<Response> {
1743 let write = git.service == GitService::ReceivePack;
1744 let get = request.method() == Method::Get;
1745 let identity = env.service("IDENTITY")?;
1746 // The repository and the caller's credentials, at once. A fetch may
1747 // go by the row as read a moment ago, for the same clone's next
1748 // request; a push always reads it. Anonymous callers cost nothing.
1749 let lookup = async {
1750 if write {
1751 self.registry.by_path(&git.path).await
1752 } else {
1753 self.registry.by_path_recent(&git.path).await
1754 }
1755 };
1756 let (found, viewer) =
1757 futures_util::future::join(lookup, git_http::viewer(&request, &identity)).await;
1758 let mut found = found?;
1759 timing.mark("repo");
1760 // A workspace alias staff set (identity's aliases.rs: `g1t` for
1761 // `flagon-io`) is answered in place, as the repository under the
1762 // workspace's slug: pushes and some clients do not follow
1763 // redirects. Everything after this sees only the workspace's slug.
1764 let aliased = match found {
1765 Some(_) => None,
1766 None => git_http::aliased(git, &identity).await?,
1767 };
1768 if let Some(aliased) = &aliased {
1769 found = if write {
1770 self.registry.by_path(&aliased.path).await?
1771 } else {
1772 self.registry.by_path_recent(&aliased.path).await?
1773 };
1774 timing.mark("alias");
1775 }
1776 let git = aliased.as_ref().unwrap_or(git);
1777 if found.is_none() {
1778 // A workspace that was renamed: git follows a redirect when it
1779 // first asks for refs, and uses the new address from then on.
1780 // A repository transferred to another workspace: the same, to
1781 // its new path. Fetches and pushes both follow either.
1782 let url = request.url()?;
1783 let (renamed, moved) = futures_util::future::join(
1784 git_http::renamed(&url, &identity),
1785 self.registry.resolve_moved(&git.path),
1786 )
1787 .await;
1788 timing.mark("moved");
1789 if let Some(location) = renamed? {
1790 return git_http::moved(&location, get);
1791 }
1792 if let Some(now) = moved?
1793 && let Some(location) = git_http::transferred(&url, &now)
1794 {
1795 return git_http::moved(&location, get);
1796 }
1797 }
1798 let viewer = viewer?;
1799 // A person who has not confirmed their email address can do
1800 // nothing with git until they do: told so, not asked to sign in.
1801 if viewer.as_ref().is_some_and(g1t_contracts::User::awaits_confirmation) {
1802 let site = request.url()?.origin().ascii_serialization();
1803 return git_http::refuse(Outcome::<()>::fail(
1804 FailureCode::Forbidden,
1805 g1t_contracts::accounts::confirm_email_first(&site),
1806 ));
1807 }
1808 // A run credential is checked against its grants, then acts as the
1809 // person it works for. See run_access.rs.
1810 let (request, viewer, audit) = match self.admit_git(request, git, viewer, found.as_ref()).await? {
1811 run_access::Admitted::Go { request, viewer, entry } => (request, viewer, entry),
1812 run_access::Admitted::Refused(response) => return Ok(response),
1813 };
1814 let mut after = AfterGit {
1815 audit,
1816 status: 0,
1817 message: None,
1818 push: None,
1819 };
1820 let repo = match self.authorize_git(&git.path, &viewer, git.service, found).await? {
1821 Outcome::Ok(repo) => repo,
1822 refused => {
1823 let response = git_http::refuse(refused)?;
1824 after.ended(response.status_code(), None);
1825 after.spawn(env, ctx);
1826 return Ok(response);
1827 }
1828 };
1829 // A pull request's working copy removed after it closed is made
1830 // again before git uses it (forks.rs).
1831 self.live(&repo).await?;
1832 timing.mark("access");
1833 // Clones check out the default branch g1t keeps, which can have
1834 // changed since the store made the repository.
1835 let default_branch = repo.fork_of.is_none().then(|| repo.default_branch.clone());
1836 let key = store_key(&repo);
1837 let scope = if write { Scope::Write } else { Scope::Read };
1838 let mut request = request;
1839 let protocol = refs_cache::protocol(request.headers().get("git-protocol")?.as_deref());
1840 // A fetch's POST is read here, to tell an `ls-refs` from a fetch of
1841 // objects; the store would have it read in full anyway.
1842 let body = if !write && !get { Some(request.bytes().await?) } else { None };
1843 // What it asks the store, for the meters (meters.rs).
1844 let call = git_ops::classify(git.service, git.endpoint, get, body.as_deref());
1845 // Answers kept from the usual store may name refs the fallback
1846 // store does not have (fallback.rs): none are used, or kept.
1847 let fallback = self.store.on_fallback(&key);
1848 // An answer that lists refs may have been kept: see refs_cache.rs.
1849 let kept_key = refs_cache::kind(git, get, protocol, body.as_deref())
1850 .filter(|_| !fallback)
1851 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1852 .map(|(kind, version)| {
1853 refs_cache::Key::new(&repo.id, version, default_branch.as_deref(), protocol, &kind)
1854 });
1855 // A fresh clone's pack may have been kept too: see pack_cache.rs.
1856 // Under the same refs version, so never across a change to them.
1857 let pack_key = self
1858 .packs
1859 .as_ref()
1860 .filter(|_| !fallback)
1861 .and_then(|_| {
1862 let encoding = request.headers().get("content-encoding").ok().flatten();
1863 pack_cache::cacheable(git, get, protocol, encoding.as_deref(), body.as_deref())
1864 })
1865 .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1866 .map(|(normalized, version)| pack_cache::Key::new(&repo.id, version, &normalized));
1867 // A kept answer and the free workspace limits, with a kept
1868 // credential looked up alongside. A kept answer goes back without
1869 // waiting for the credential, which it does not need.
1870 let ((answer, pack, limited), kept_access) = {
1871 let shared = self.shared.as_deref();
1872 let answer_and_limits = std::pin::pin!(futures_util::future::join3(
1873 async {
1874 match &kept_key {
1875 Some(kept_key) => refs_cache::get(shared, kept_key).await,
1876 None => None,
1877 }
1878 },
1879 async {
1880 match (&pack_key, self.packs.as_deref()) {
1881 (Some(pack_key), Some(packs)) => pack_cache::get(packs, pack_key).await,
1882 _ => None,
1883 }
1884 },
1885 self.git_limits(call, git, &repo, env),
1886 ));
1887 let kept_access = std::pin::pin!(self.store.kept_access(&key, scope));
1888 match futures_util::future::select(answer_and_limits, kept_access).await {
1889 futures_util::future::Either::Left((first, kept_access)) => {
1890 let answered = first.0.is_some() || first.1.is_some() || matches!(first.2, Ok(Some(_)) | Err(_));
1891 (first, if answered { None } else { kept_access.await })
1892 }
1893 futures_util::future::Either::Right((kept_access, first)) => (first.await, kept_access),
1894 }
1895 };
1896 timing.mark("kept");
1897 // A kept pack first: it never reaches the store, so it is never an
1898 // operation, and a free workspace past its operation cap still gets
1899 // it. (The other limits are a push's, and a pack is only a fetch.)
1900 if let Some(kept) = pack {
1901 timing.note("pack", "hit");
1902 let sent = body.as_ref().map_or(0, |body| body.len() as u64);
1903 meters::record(pack_cache::HIT, &key, sent, kept.size);
1904 after.ended(200, None);
1905 after.spawn(env, ctx);
1906 return kept.response();
1907 }
1908 if let Some((response, status, message)) = limited? {
1909 after.ended(status, Some(message.to_owned()));
1910 after.spawn(env, ctx);
1911 return Ok(response);
1912 }
1913 if pack_key.is_some() {
1914 timing.note("pack", "miss");
1915 }
1916 if let (Some((entry, found)), Some(kept_key)) = (answer, &kept_key) {
1917 timing.note("refs", found.as_str());
1918 if found == refs_cache::Found::Shared {
1919 let (kept_key, entry) = (kept_key.clone(), entry.clone());
1920 ctx.wait_until(async move { refs_cache::keep_in_colo(&kept_key, &entry).await });
1921 }
1922 // Never reached the store: never an operation.
1923 meters::record(call.cached_meter(), &key, 0, entry.body.len() as u64);
1924 after.ended(200, None);
1925 after.spawn(env, ctx);
1926 return entry.response();
1927 }
1928 if kept_key.is_some() {
1929 timing.note("refs", "miss");
1930 }
1931 // The store's credential: one made a moment ago, here or in another
1932 // isolate (see store.rs), or a new one.
1933 let access = match kept_access {
1934 Some((access, from)) => {
1935 timing.note("cred", from.as_str());
1936 access
1937 }
1938 None => {
1939 let access = self.store.mint_access(&key, scope).await?;
1940 timing.mark("mint");
1941 timing.note("cred", "mint");
1942 access
1943 }
1944 };
1945 // Should the store turn a kept credential down, a fetch's first
1946 // request is tried again with a new one; the requests after it then
1947 // have that one too.
1948 let again = if get { Some(request.clone()?) } else { None };
1949 // Push protection: a push that adds a secret is refused. See secret_scan.rs.
1950 let scan = async |body: &[u8]| self.protect(&repo, viewer.as_ref(), body).await;
1951 // Rulesets: what the rules of the branches and tags it changes
1952 // refuse is declined, saying which rule and why (rules.rs).
1953 let rules = async |head: &[u8], whole: bool| self.check_push(&repo, viewer.as_ref(), head, whole).await;
1954 // What a push may bring (pack_limits.rs): the repository's size is
1955 // its own and its pull requests' working copies'.
1956 let limits = if write && !get {
1957 git_http::PushLimits {
1958 held: self.held(&repo).await,
1959 repo_limit: self.repo_limit,
1960 large: self.large_pushes,
1961 ..git_http::PushLimits::default()
1962 }
1963 } else {
1964 git_http::PushLimits::default()
1965 };
1966 let mut outcome = git_http::forward(
1967 request,
1968 body,
1969 git,
1970 &access,
1971 rules,
1972 default_branch.as_deref(),
1973 limits,
1974 scan,
1975 )
1976 .await?;
1977 let turned_down = matches!(
1978 &outcome,
1979 git_http::Push::Forwarded(forwarded) if matches!(forwarded.response.status_code(), 401 | 403)
1980 );
1981 if turned_down {
1982 self.store.forget_access(&key).await;
1983 if let Some(again) = again {
1984 let access = self.store.mint_access(&key, scope).await?;
1985 let nothing = async |_: &[u8]| Ok(None);
1986 outcome = git_http::forward(
1987 again,
1988 None,
1989 git,
1990 &access,
1991 async |_: &[u8], _: bool| Ok(None),
1992 default_branch.as_deref(),
1993 git_http::PushLimits::default(),
1994 nothing,
1995 )
1996 .await?;
1997 }
1998 }
1999 let forwarded =
2000 match outcome {
2001 git_http::Push::Forwarded(forwarded) => forwarded,
2002 git_http::Push::Refused(response) => {
2003 after.ended(403, Some("The push was declined by rules.".to_owned()));
2004 after.spawn(env, ctx);
2005 return Ok(response);
2006 }
2007 git_http::Push::Blocked(response) => {
2008 after.ended(403, Some("The push adds a secret.".to_owned()));
2009 after.spawn(env, ctx);
2010 return Ok(response);
2011 }
2012 git_http::Push::Declined(response, reason) => {
2013 after.ended(403, Some(format!("The push was declined: {reason}.")));
2014 after.spawn(env, ctx);
2015 return Ok(response);
2016 }
2017 };
2018 if forwarded.from_store {
2019 let received = forwarded
2020 .response
2021 .headers()
2022 .get("content-length")?
2023 .and_then(|length| length.parse().ok())
2024 .unwrap_or(0);
2025 meters::record(call.meter(), &key, forwarded.sent, received);
2026 }
2027 timing.mark("store");
2028 let mut response = forwarded.response;
2029 let status = response.status_code();
2030 if write && !get {
2031 // A push: the store has moved its refs once it has answered in
2032 // full, so the answer is read before the change is recorded, and
2033 // only then goes back. Whoever fetches after it sees the push.
2034 let headers = response.headers().clone();
2035 headers.delete("content-length")?;
2036 let report = response.bytes().await?;
2037 self.refs_moved(&repo.id).await;
2038 timing.mark("refs");
2039 response = Response::from_bytes(report)?.with_headers(headers).with_status(status);
2040 } else if let (Some(kept_key), 200) = (&kept_key, status) {
2041 // A miss: this answer is kept for the next to ask.
2042 let headers = response.headers().clone();
2043 headers.delete("content-length")?;
2044 let body = response.bytes().await?;
2045 if let Some(content_type) = headers.get("content-type")? {
2046 let entry = refs_cache::Entry { content_type, body: body.clone() };
2047 if entry.keepable() {
2048 let shared = self.shared.clone();
2049 let kept_key = kept_key.clone();
2050 ctx.wait_until(async move { refs_cache::keep(shared.as_deref(), &kept_key, &entry).await });
2051 }
2052 }
2053 response = Response::from_bytes(body)?.with_headers(headers).with_status(status);
2054 } else if let (Some(pack_key), Some(packs), true) = (&pack_key, &self.packs, forwarded.from_store) {
2055 // A fresh clone the bucket did not have: counted, and its pack
2056 // kept as it streams to git, when it is a whole one.
2057 meters::record(pack_cache::MISS, &key, forwarded.sent, 0);
2058 if status == 200 {
2059 let store_key = key.clone();
2060 let measured = Box::new(move |bytes: u64| meters::record_bytes(pack_cache::MISS, &store_key, 0, bytes));
2061 let (teed, filling) = pack_cache::tee(response, packs.clone(), pack_key, measured)?;
2062 response = teed;
2063 if let Some(filling) = filling {
2064 let pack_key = pack_key.clone();
2065 ctx.wait_until(async move {
2066 let filled = filling.await;
2067 if !matches!(filled, pack_cache::Filled::Kept { .. } | pack_cache::Filled::Abandoned) {
2068 worker::console_warn!("pack {} not kept: {filled:?}", pack_key.as_str());
2069 }
2070 });
2071 }
2072 }
2073 }
2074 after.ended(status, None);
2075 if status == 200 && (forwarded.pack_bytes > 0 || !forwarded.pushed.is_empty()) {
2076 after.push = Some(PushDone {
2077 repo,
2078 pushed: forwarded.pushed,
2079 pack_bytes: forwarded.pack_bytes,
2080 caused_by_job: viewer.as_ref().and_then(g1t_contracts::events::job_run_of).map(str::to_owned),
2081 actor: viewer.map(|user: User| user.id),
2082 unscanned: forwarded.unscanned,
2083 });
2084 }
2085 after.spawn(env, ctx);
2086 Ok(response)
2087 }
2088
2089 /// The answer for a request a free workspace's limits stop, or a push
2090 /// to a full repository, with its status and reason for the audit log;
2091 /// `None` to go on.
2092 ///
2093 /// A clone, fetch or push is a git operation, which the git store
2094 /// charges g1t for: counted for billing once the answer has gone back
2095 /// (meters.rs), and a free workspace far past its share is slowed down
2096 /// rather than charged (see git_ops.rs). Whether it is past it is
2097 /// decided from counts this isolate already holds: the database is not
2098 /// asked on the way. A free workspace is never charged for private
2099 /// storage: once its private repositories hold the free amount, pushes
2100 /// to them stop, checked when a push begins so that git shows the
2101 /// reason. So do pushes to a repository at the store's size limit.
2102 async fn git_limits(
2103 &self,
2104 call: git_ops::GitCall,
2105 git: &git_http::GitRequest,
2106 repo: &Repo,
2107 env: &Env,
2108 ) -> Result<Option<(Response, u16, &'static str)>> {
2109 let namespace = git.path.namespace.to_lowercase();
2110 if meters::mapping_now().billable(call.meter()) > 0.0 {
2111 let now = now_ms();
2112 let hour = git_ops::hour_key(&rfc3339(now));
2113 let limits = git_ops::Limits::from_env(env);
2114 if let Some((month, hour_ops)) = git_ops::standing(&namespace, &hour, now)
2115 && git_ops::slow_down(month + 1, hour_ops + 1, limits.free_cap, limits.hourly)
2116 && git_ops::is_free_kept(env.service("BILLING").ok().as_ref(), &namespace).await
2117 {
2118 return Ok(Some((
2119 git_ops::too_many(&namespace, limits.free_cap, limits.hourly)?,
2120 429,
2121 "Too many git operations this hour.",
2122 )));
2123 }
2124 }
2125 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" {
2126 let held = self.held(repo).await;
2127 if held >= self.repo_limit {
2128 let message = format!(
2129 "{}/{} 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",
2130 repo.namespace,
2131 repo.name,
2132 pack_limits::megabytes(held)
2133 );
2134 return Ok(Some((Response::error(message, 403)?, 403, "The repository is full.")));
2135 }
2136 }
2137 if git.service == GitService::ReceivePack && git.endpoint == "info/refs" && repo.is_private {
2138 let free = git_ops::free_private_bytes(env);
2139 let held = self.registry.private_bytes(&namespace).await.unwrap_or(0);
2140 if git_ops::storage_full(held, free)
2141 && git_ops::is_free(env.service("BILLING").ok().as_ref(), &namespace).await
2142 {
2143 return Ok(Some((
2144 git_ops::storage_full_response(&namespace, held, free)?,
2145 403,
2146 "Free private storage is full.",
2147 )));
2148 }
2149 }
2150 Ok(None)
2151 }
2152
2153 /// What a repository and its pull requests' working copies hold, as
2154 /// g1t counts it: read for a push's first request, kept a minute for
2155 /// the rest of it.
2156 async fn held(&self, repo: &Repo) -> u64 {
2157 let root = repo.fork_of.clone().unwrap_or_else(|| repo.id.clone());
2158 let now = now_ms();
2159 if let Some(held) = HELD.with(|held| held.borrow().get(&root, now)) {
2160 return held;
2161 }
2162 let held = self.registry.stored_bytes(&root).await.unwrap_or(0).max(0) as u64;
2163 HELD.with(|kept| kept.borrow_mut().put(root, held, now));
2164 held
2165 }
2166
2167 /// What a push changed, recorded once git has its answer.
2168 async fn record_push(&self, push: PushDone) -> Result<()> {
2169 let PushDone {
2170 repo,
2171 pushed,
2172 pack_bytes,
2173 actor,
2174 unscanned,
2175 caused_by_job,
2176 } = push;
2177 // What the push stored, for billing's storage meter. A failure only
2178 // leaves the count short.
2179 if pack_bytes > 0
2180 && let Err(error) = self.registry.add_stored_bytes(&repo, pack_bytes).await
2181 {
2182 worker::console_error!("stored bytes for {} not counted: {error}", repo.name);
2183 }
2184 if pushed.is_empty() {
2185 return Ok(());
2186 }
2187 // Artifacts' own push notifications are per repository, which does
2188 // not fit a repo per pull request, so the front end reports pushes
2189 // itself: one event for each branch that moved.
2190 let stored = self.store.open(&store_key(&repo)).await?;
2191 let announced = announced_refs(&pushed, &repo.default_branch);
2192 if announced.len() < pushed.len() {
2193 worker::console_log!(
2194 "push to {}: {} of {} refs announced",
2195 repo.name,
2196 announced.len(),
2197 pushed.len()
2198 );
2199 }
2200 for pushed in announced {
2201 // The store can refuse one ref and accept another, so each
2202 // branch is checked against where it actually is. A tag the
2203 // store cannot read back is taken as pushed.
2204 let moved = match pushed.branch() {
2205 Some(branch) => stored
2206 .log(branch, 1)
2207 .await?
2208 .first()
2209 .is_some_and(|commit| commit.hash == pushed.after),
2210 None => stored.log(&pushed.git_ref, 1).await.map_or(true, |head| {
2211 head.first().is_none_or(|commit| commit.hash == pushed.after)
2212 }),
2213 };
2214 if moved {
2215 self.publish_git_push(
2216 &repo,
2217 &pushed.git_ref,
2218 pushed.before.as_deref(),
2219 &pushed.after,
2220 actor.clone(),
2221 unscanned,
2222 caused_by_job.clone(),
2223 )
2224 .await?;
2225 }
2226 }
2227 Ok(())
2228 }
2229}
2230
2231/// The most tags one push announces: a push of more (`git push --tags`
2232/// into a new repository, say) announces none of them, so it starts no
2233/// workflows, mirror syncs or package reads, one per tag.
2234const MAX_PUSH_TAG_EVENTS: usize = 3;
2235/// The most branches one push announces. A push of more announces only
2236/// the default branch, if it moved, which the rest of g1t reads from.
2237const MAX_PUSH_BRANCH_EVENTS: usize = 1000;
2238
2239/// The refs of a push that are announced with a `git.push` event each.
2240/// Every ref is stored whatever this says; only the events are capped.
2241fn announced_refs<'a>(pushed: &'a [git_http::Pushed], default_branch: &str) -> Vec<&'a git_http::Pushed> {
2242 let (branches, tags): (Vec<&git_http::Pushed>, Vec<&git_http::Pushed>) =
2243 pushed.iter().partition(|pushed| pushed.branch().is_some());
2244 let mut announced = if branches.len() > MAX_PUSH_BRANCH_EVENTS {
2245 branches.into_iter().filter(|pushed| pushed.branch() == Some(default_branch)).collect()
2246 } else {
2247 branches
2248 };
2249 if tags.len() <= MAX_PUSH_TAG_EVENTS {
2250 announced.extend(tags);
2251 }
2252 announced
2253}
2254
2255/// A push the store accepted, to be recorded once git has its answer.
2256struct PushDone {
2257 repo: Repo,
2258 pushed: Vec<git_http::Pushed>,
2259 pack_bytes: u64,
2260 actor: Option<String>,
2261 /// Too large to scan for secrets before it was stored.
2262 unscanned: bool,
2263 /// The run whose job's token pushed, if one did: its push starts no
2264 /// workflows.
2265 caused_by_job: Option<String>,
2266}
2267
2268/// What a git request leaves for after its answer: its audit entry, with
2269/// how the request ended, and what a push changed.
2270struct AfterGit {
2271 audit: Option<Box<g1t_contracts::audit::NewAuditEntry>>,
2272 status: u16,
2273 message: Option<String>,
2274 push: Option<PushDone>,
2275}
2276
2277impl AfterGit {
2278 fn ended(&mut self, status: u16, message: Option<String>) {
2279 self.status = status;
2280 self.message = message;
2281 }
2282
2283 /// Does the work once the response is on its way. A failure is logged:
2284 /// git has already been told how its request went.
2285 fn spawn(self, env: &Env, ctx: &Context) {
2286 if self.audit.is_none() && self.push.is_none() {
2287 return;
2288 }
2289 let env = env.clone();
2290 ctx.wait_until(async move {
2291 let repos = match service(&env) {
2292 Ok(repos) => repos,
2293 Err(error) => {
2294 worker::console_error!("git request not recorded: {error}");
2295 return;
2296 }
2297 };
2298 repos.finish_git(self.audit, self.status, self.message).await;
2299 if let Some(push) = self.push
2300 && let Err(error) = repos.record_push(push).await
2301 {
2302 worker::console_error!("push not recorded: {error}");
2303 }
2304 });
2305 }
2306}
2307
2308fn service(env: &Env) -> Result<Repos<ArtifactsStore>> {
2309 let shared = shared::Shared::from_env(env).map(Rc::new);
2310 Ok(Repos {
2311 registry: Registry { db: env.d1("DB")? },
2312 store: ArtifactsStore::new(env, shared.clone())?,
2313 shared,
2314 packs: pack_cache::Packs::from_env(env).map(Rc::new),
2315 events: env.service("EVENTS")?,
2316 security: env.service("SECURITY").ok(),
2317 billing: env.service("BILLING").ok(),
2318 identity: env.service("IDENTITY").ok(),
2319 work: env.service("WORK").ok(),
2320 free_private_bytes: git_ops::free_private_bytes(env),
2321 fork_days: forks::retention_days(env),
2322 repo_limit: env
2323 .var("REPO_STORAGE_LIMIT_BYTES")
2324 .ok()
2325 .and_then(|value| value.to_string().parse().ok())
2326 .unwrap_or(pack_limits::DEFAULT_REPO_LIMIT_BYTES),
2327 large_pushes: git_http::LargePushes::from_var(env.var("LARGE_PUSHES").ok().map(|value| value.to_string()).as_deref()),
2328 placement: shards::Placement::from_vars(
2329 env.var("ARTIFACTS_NEW_REPOS").ok().map(|value| value.to_string()).as_deref(),
2330 env.var("ARTIFACTS_EU_NAMESPACE").ok().map(|value| value.to_string()).as_deref(),
2331 ),
2332 limits: shards::limits(env.var("ARTIFACTS_NAMESPACE_LIMITS").ok().map(|value| value.to_string()).as_deref()),
2333 })
2334}
2335
2336/// Writes what this isolate metered once the answer has gone back, every
2337/// few seconds at most: now, or once it is due, waiting in this request's
2338/// `wait_until` so nothing counted is left for a request that may never
2339/// come (meters.rs).
2340fn flush_later(env: &Env, ctx: &Context) {
2341 let Some(wait) = meters::plan_flush() else {
2342 return;
2343 };
2344 if let Ok(db) = env.d1("DB") {
2345 ctx.wait_until(async move { meters::flush_after(&db, wait).await });
2346 }
2347}
2348
2349/// `/backups/<job id>/parts/<number>`: the job and the part's number.
2350fn backup_part_path(path: &str) -> Option<(String, u16)> {
2351 let rest = path.strip_prefix("/backups/")?;
2352 let (job, number) = rest.split_once("/parts/")?;
2353 let number = number.parse::<u16>().ok()?;
2354 (!job.is_empty() && !job.contains('/')).then(|| (job.to_owned(), number))
2355}
2356
2357fn backups_off<T>() -> Outcome<T> {
2358 Outcome::fail(FailureCode::Conflict, "Backups are off on this installation: it has no storage for them.")
2359}
2360
2361/// One part of a backup's bundle, with the job's token in its header.
2362async fn backup_part(request: &mut Request, env: &Env, repos: &Repos<ArtifactsStore>, job_id: String, number: u16) -> Result<Response> {
2363 let Some(blobs) = backups::storage(env) else {
2364 return reply(&backups_off::<()>());
2365 };
2366 let token = request.headers().get(g1t_contracts::backups::TOKEN_HEADER)?.unwrap_or_default();
2367 let bytes = request.bytes().await?;
2368 let job = g1t_contracts::backups::BackupJobArgs { job_id, token };
2369 reply(&backups::part(&repos.registry.db, &blobs, &job, number, bytes).await?)
2370}
2371
2372#[cfg(test)]
2373mod backup_path_tests {
2374 use super::backup_part_path;
2375
2376 #[test]
2377 fn a_part_is_named_by_its_job_and_number() {
2378 assert_eq!(backup_part_path("/backups/bkp_1/parts/3"), Some(("bkp_1".to_owned(), 3)));
2379 assert_eq!(backup_part_path("/backups/bkp_1/parts/x"), None);
2380 assert_eq!(backup_part_path("/backups//parts/1"), None);
2381 assert_eq!(backup_part_path("/acme/rocket.git/info/refs"), None);
2382 }
2383}
2384
2385/// Read methods whose answer is an `Outcome`: when the git store is busy,
2386/// the site is told so in words instead of failing the page.
2387const OUTCOME_READS: [&str; 7] = ["tree", "blob", "log", "branches", "blame", "compare", "branch_drift"];
2388
2389#[event(fetch)]
2390async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
2391 let mut repos = service(&env)?;
2392 // A part of a backup's bundle, as the API passes it on from the
2393 // sandbox: bytes, not JSON (backups.rs).
2394 if request.method() == Method::Put
2395 && let Some((job_id, number)) = backup_part_path(&request.path())
2396 {
2397 let answered = backup_part(&mut request, &env, &repos, job_id, number).await;
2398 flush_later(&env, &ctx);
2399 return answered;
2400 }
2401 let Some(method) = rpc_method(&request) else {
2402 let answered = repos.git_http(request, &env, &ctx).await;
2403 flush_later(&env, &ctx);
2404 return answered;
2405 };
2406 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2407 // Git over HTTPS above always reads the primary.
2408 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2409 repos.registry.db = db;
2410 let body: serde_json::Value = request.json().await?;
2411
2412 let answered = async { match method.as_str() {
2413 "get" => reply(&repos.get(args(body)?).await?),
2414 "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
2415 "readable" => {
2416 let a: ReadableArgs = args(body)?;
2417 reply(&repos.registry.readable(&a.ids, &a.viewer).await?)
2418 }
2419 "public_namespaces" => {
2420 let a: PublicNamespacesArgs = args(body)?;
2421 reply(&repos.registry.public_namespaces(&a.owner_id).await?)
2422 }
2423 "path_by_id" => {
2424 let a: PathByIdArgs = args(body)?;
2425 reply(
2426 &repos
2427 .registry
2428 .by_id(&a.id)
2429 .await?
2430 .filter(|repo| repo.fork_of.is_none())
2431 .map(|repo| RepoPath {
2432 namespace: repo.namespace,
2433 name: repo.name,
2434 }),
2435 )
2436 }
2437 "list" => {
2438 let a: ListArgs = args(body)?;
2439 reply(
2440 &repos
2441 .registry
2442 .list(
2443 &a.viewer,
2444 a.query.as_deref(),
2445 a.namespace.as_deref(),
2446 a.member_only,
2447 )
2448 .await?,
2449 )
2450 }
2451 "create" => reply(&repos.create(args(body)?).await?),
2452 // Services only: a GitHub mirror catching up, or pushing out.
2453 "mirror" => reply(&repos.mirror(args(body)?).await?),
2454 "transfer" => reply(&repos.transfer(args(body)?).await?),
2455 // A repository's lifecycle: see lifecycle.rs.
2456 "delete" => reply(&repos.delete(args(body)?).await?),
2457 "deleted" => reply(&repos.deleted(args(body)?).await?),
2458 "restore" => reply(&repos.restore(args(body)?).await?),
2459 "purge" => reply(&repos.purge(args(body)?).await?),
2460 "purge_due" => reply(&repos.purge_due(args(body)?).await?),
2461 "rename" => reply(&repos.rename(args(body)?).await?),
2462 "archive" => reply(&repos.archive(args(body)?).await?),
2463 "set_visibility" => reply(&repos.set_visibility(args(body)?).await?),
2464 "set_default_branch" => reply(&repos.set_default_branch(args(body)?).await?),
2465 "rename_branch" => reply(&repos.rename_branch(args(body)?).await?),
2466 "resolve_branch" => reply(&repos.resolve_branch(args(body)?).await?),
2467 "status_by_id" => reply(&repos.status_by_id(args(body)?).await?),
2468 "resolve_path" => {
2469 let a: ResolvePathArgs = args(body)?;
2470 reply(&repos.registry.resolve_moved(&a.path).await?)
2471 }
2472 "namespace_count" => {
2473 let a: NamespaceCountArgs = args(body)?;
2474 reply(&repos.registry.count_in(&a.namespace).await?)
2475 }
2476 "update" => reply(&repos.update(args(body)?).await?),
2477 "tree" => reply(&repos.tree(args(body)?).await?),
2478 "blob" => reply(&repos.blob(args(body)?).await?),
2479 "log" => reply(&repos.log(args(body)?).await?),
2480 "blame" => reply(&repos.blame(args(body)?).await?),
2481 "fork_for_pull" => reply(&repos.fork_for_pull(args(body)?).await?),
2482 "git_access" => reply(&repos.git_access(args(body)?).await?),
2483 "branches" => reply(&repos.branches(args(body)?).await?),
2484 "last_commits" => reply(&repos.last_commits(args(body)?).await?),
2485 "branch_drift" => reply(&repos.branch_drift(args(body)?).await?),
2486 "tags" => reply(&repos.tags(args(body)?).await?),
2487 // The About: what the Files page shows beside the files (about.rs).
2488 // What is kept behind the head is worked out again after the answer.
2489 "about" => {
2490 let (answer, refresh) = repos.about(args(body)?).await?;
2491 about::refresh_later(&env, &ctx, refresh);
2492 reply(&answer)
2493 }
2494 "languages" => {
2495 let (answer, refresh) = repos.languages(args(body)?).await?;
2496 about::refresh_later(&env, &ctx, refresh);
2497 reply(&answer)
2498 }
2499 "contributors" => {
2500 let (answer, refresh) = repos.contributors(args(body)?).await?;
2501 about::refresh_later(&env, &ctx, refresh);
2502 reply(&answer)
2503 }
2504 "license" => {
2505 let (answer, refresh) = repos.license(args(body)?).await?;
2506 about::refresh_later(&env, &ctx, refresh);
2507 reply(&answer)
2508 }
2509 "stars" => reply(&repos.stars(args(body)?).await?),
2510 "star" => reply(&repos.star(args(body)?).await?),
2511 "stargazers" => reply(&repos.stargazers(args(body)?).await?),
2512 "starred" => reply(&repos.starred(args(body)?).await?),
2513 "releases" => reply(&repos.releases(args(body)?).await?),
2514 "release" => reply(&repos.release(args(body)?).await?),
2515 "create_release" => reply(&repos.create_release(args(body)?).await?),
2516 "update_release" => reply(&repos.update_release(args(body)?).await?),
2517 "delete_release" => reply(&repos.delete_release(args(body)?).await?),
2518 "head" => reply(&repos.head(args(body)?).await?),
2519 "behind" => reply(&repos.behind(args(body)?).await?),
2520 "divergence" => reply(&repos.divergence(args(body)?).await?),
2521 "land" => reply(&repos.land(args(body)?).await?),
2522 "update_pull_branch" => reply(&repos.update_pull_branch(args(body)?).await?),
2523 "delete_branch" => reply(&repos.delete_branch(args(body)?).await?),
2524 "commit_file" => reply(&repos.commit_file(args(body)?).await?),
2525 "compare" => reply(&repos.compare(args(body)?).await?),
2526 // Services only: a pull request's commits, as rules look at them (rules.rs).
2527 "inspect_commits" => reply(&repos.inspect_commits(args(body)?).await?),
2528 "scan_history" => reply(&repos.scan_history(args(body)?).await?),
2529 "find_lockfiles" => reply(&repos.find_lockfiles(args(body)?).await?),
2530 "match_pattern" => reply(&repos.match_pattern(args(body)?).await?),
2531 "check_secret" => reply(&repos.check_secret(args(body)?).await?),
2532 "list_files" => reply(&repos.list_files(args(body)?).await?),
2533 "changed_files" => reply(&repos.changed_files(args(body)?).await?),
2534 "read_blobs" => reply(&repos.read_blobs(args(body)?).await?),
2535 // Services only: what the Composer registry builds packages from.
2536 "refs" => reply(&repos.refs_of(args(body)?).await?),
2537 "raw_file" => reply(&repos.raw_file(args(body)?).await?),
2538 "raw_blobs" => reply(&repos.raw_blobs(args(body)?).await?),
2539 "visibility" => {
2540 let a: g1t_contracts::repos::VisibilityArgs = args(body)?;
2541 reply(&repos.registry.visibility(&a.paths).await?)
2542 }
2543 "storage" => reply(&repos.registry.storage().await?),
2544 "git_operations" => {
2545 let a: GitOperationsArgs = args(body)?;
2546 reply(&git_ops::totals(&repos.registry.db, &a.month, a.since.as_deref(), a.namespace.as_deref().map(str::to_lowercase).as_deref()).await?)
2547 }
2548 // Identity, once: who created each repository (members.rs there).
2549 "repo_creators" => {
2550 let a: AllIdsArgs = args(body)?;
2551 let limit = a.limit.clamp(1, 500);
2552 let repos = repos.registry.creators_after(a.after.as_deref(), limit).await?;
2553 let next = (repos.len() == limit as usize).then(|| repos.last().map(|repo| repo.id.clone())).flatten();
2554 reply(&CreatorPage { repos, next })
2555 }
2556 "all_ids" => {
2557 let a: AllIdsArgs = args(body)?;
2558 let limit = a.limit.clamp(1, 500);
2559 let ids = repos.registry.ids_after(a.after.as_deref(), limit).await?;
2560 let next = (ids.len() == limit as usize).then(|| ids.last().cloned()).flatten();
2561 reply(&IdPage { ids, next })
2562 }
2563 // The raw meters of the git store, for reconciling with Cloudflare
2564 // (meters.rs, scripts/ops/artifacts-usage.mjs).
2565 "artifacts_usage" => {
2566 let a: meters::UsageArgs = args(body)?;
2567 reply(&meters::usage(&repos.registry.db, &a).await?)
2568 }
2569 "operation_mapping" => reply(&meters::read_mapping(&repos.registry.db).await?),
2570 // Billing: the workspace each pull request's working copy is counted
2571 // for, so Cloudflare's own count of `pulls--<id>` shares out too.
2572 "pull_owners" => {
2573 #[derive(serde::Deserialize)]
2574 struct PullOwnersArgs {
2575 pulls: Vec<String>,
2576 }
2577 let a: PullOwnersArgs = args(body)?;
2578 let pulls: Vec<String> = a.pulls.into_iter().take(500).collect();
2579 reply(&serde_json::json!({ "owners": meters::pull_owners(&repos.registry.db, &pulls).await? }))
2580 }
2581 // Services only: which meters are operations, changed without a deploy.
2582 "set_operation_mapping" => {
2583 let row: meters::MappingRow = args(body)?;
2584 meters::set_mapping(&repos.registry.db, &row, &rfc3339(now_ms())).await?;
2585 reply(&meters::read_mapping(&repos.registry.db).await?)
2586 }
2587 // Backups (backups.rs): the runner's sweep claims queued ones, and
2588 // each sandbox, through the API, asks for its job and says how it went.
2589 "claim_backups" => {
2590 let a: g1t_contracts::backups::ClaimBackupsArgs = args(body)?;
2591 let blobs = backups::storage(&env);
2592 reply(&backups::claim(&repos.registry.db, blobs.as_ref(), &a, now_ms()).await?)
2593 }
2594 "backup_spec" => match backups::storage(&env) {
2595 Some(blobs) => {
2596 let a: g1t_contracts::backups::BackupJobArgs = args(body)?;
2597 let every = backups::Settings::from_env(&env).full_every;
2598 reply(&backups::spec(&repos.registry, &blobs, &repos.store, &a, every, now_ms()).await?)
2599 }
2600 None => reply(&backups_off::<bool>()),
2601 },
2602 "backup_complete" => match backups::storage(&env) {
2603 Some(blobs) => reply(&backups::complete(&repos.registry, &blobs, &args(body)?, now_ms()).await?),
2604 None => reply(&backups_off::<bool>()),
2605 },
2606 "backup_fail" => match backups::storage(&env) {
2607 Some(blobs) => reply(&backups::fail(&repos.registry.db, &blobs, &args(body)?).await?),
2608 None => reply(&backups_off::<bool>()),
2609 },
2610 // How the git store has been answering, for the status page.
2611 "store_health" => {
2612 let a: meters::HealthArgs = args(body)?;
2613 reply(&meters::health(&repos.registry.db, &a).await?)
2614 }
2615 // Where repositories may be kept, for a workspace's settings.
2616 "storage_options" => reply(&repos.storage_options()),
2617 // Services and operators only: how each namespace stands, and
2618 // moving a repository between them (namespaces.rs, moves.rs).
2619 "namespaces" => reply(&repos.standings().await?),
2620 "move_repository" => reply(&repos.move_repository(args(body)?).await?),
2621 "repository_moves" => {
2622 let a: moves::ListMovesArgs = args(body)?;
2623 reply(&repos.registry.moves(a.limit.unwrap_or(50)).await?)
2624 }
2625 _ => Response::error("Unknown method", 404),
2626 } }
2627 .await;
2628 // The git store is busy: said in words, with when to try again.
2629 let answered = match answered {
2630 Err(error) => match resilience::busy(&error.to_string()) {
2631 Some(busy) if OUTCOME_READS.contains(&method.as_str()) => {
2632 reply(&Outcome::<()>::fail(FailureCode::Conflict, busy.message().trim()))
2633 }
2634 Some(busy) => {
2635 let response = Response::error(busy.message(), 503)?;
2636 response.headers().set("retry-after", &busy.retry_after.to_string())?;
2637 Ok(response)
2638 }
2639 None => Err(error),
2640 },
2641 answered => answered,
2642 };
2643 flush_later(&env, &ctx);
2644 served.finish(answered)
2645}
2646
2647/// The nightly cron in wrangler.jsonc: tonight's backups are queued.
2648const BACKUP_CRON: &str = "53 2 * * *";
2649
2650/// The hourly sweep: deleted repositories whose time to be restored has
2651/// passed are purged. See lifecycle.rs. And, at [`BACKUP_CRON`], the
2652/// repositories whose refs moved are queued for a backup (backups.rs).
2653#[event(scheduled)]
2654async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
2655 let repos = match service(&env) {
2656 Ok(repos) => repos,
2657 Err(error) => {
2658 worker::console_error!("repos: the sweep could not start: {error}");
2659 return;
2660 }
2661 };
2662 if event.cron() == BACKUP_CRON {
2663 let Some(blobs) = backups::storage(&env) else { return };
2664 match backups::nightly(&repos.registry.db, &blobs, backups::Settings::from_env(&env), now_ms()).await {
2665 Ok(night) => worker::console_log!("repos: queued {} backups, removed {} of purged repositories", night.queued, night.pruned),
2666 Err(error) => worker::console_error!("repos: backups could not be queued: {error}"),
2667 }
2668 return;
2669 }
2670 match repos.purge_due(PurgeDueArgs::default()).await {
2671 Ok(0) => {}
2672 Ok(count) => worker::console_log!("repos: purged {count} deleted repositories"),
2673 Err(error) => worker::console_error!("repos: the purge sweep failed: {error}"),
2674 }
2675 // Pull requests' working copies whose time has come (forks.rs).
2676 match repos.retire_due().await {
2677 Ok(0) => {}
2678 Ok(count) => worker::console_log!("repos: removed {count} pull request working copies"),
2679 Err(error) => worker::console_error!("repos: the working copy sweep failed: {error}"),
2680 }
2681 // Repositories moving between namespaces, and old copies (moves.rs).
2682 match repos.run_moves().await {
2683 Ok(0) => {}
2684 Ok(count) => worker::console_log!("repos: moved {count} repositories between namespaces"),
2685 Err(error) => worker::console_error!("repos: the move sweep failed: {error}"),
2686 }
2687 meters::flush(&repos.registry.db).await;
2688}
2689
2690/// Events from the bus. A workspace's rename: its repositories move to the
2691/// workspace's current slug, asked of identity by id, so a repeated or late
2692/// delivery lands in the same place; their git store keys stay as they
2693/// were. A workspace's deletion: its repositories are deleted with it,
2694/// restored with it, or purged with it.
2695#[event(queue)]
2696async fn queue(batch: MessageBatch<Event>, env: Env, ctx: Context) -> Result<()> {
2697 let registry = Registry { db: env.d1("DB")? };
2698 let identity = env.service("IDENTITY")?;
2699 let handled = handle_events(&batch, &env, &registry, &identity).await;
2700 flush_later(&env, &ctx);
2701 handled
2702}
2703
2704/// Each event is acknowledged or retried on its own, so one that fails is
2705/// tried again without the others before and after it running twice.
2706async fn handle_events(batch: &MessageBatch<Event>, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
2707 for message in batch.messages()? {
2708 match handle_event(message.body(), env, registry, identity).await {
2709 Ok(()) => message.ack(),
2710 Err(error) => {
2711 worker::console_error!("repos: event {} failed: {error}", message.body().id);
2712 message.retry();
2713 }
2714 }
2715 }
2716 Ok(())
2717}
2718
2719async fn handle_event(event: &Event, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
2720 // A pull request merged, closed or reopened: its working copy is
2721 // kept or let go (forks.rs).
2722 if let Some(change) = forks::pull_change(&event.kind) {
2723 let Some(pull_id) = forks::pull_id_of(&event.data) else {
2724 worker::console_error!("{} {} names no pull request", event.kind, event.id);
2725 return Ok(());
2726 };
2727 let repos = service(env)?;
2728 match change {
2729 forks::PullChange::Settled => repos.pull_settled(&pull_id).await?,
2730 forks::PullChange::Reopened => repos.pull_reopened(&pull_id).await?,
2731 }
2732 return Ok(());
2733 }
2734 // A workspace deleted, restored or purged: its repositories go with
2735 // it, come back with it, or are purged with it (lifecycle.rs).
2736 if event.kind == "workspace.deleting" {
2737 match serde_json::from_value::<WorkspaceDeleting>(event.data.clone()) {
2738 Ok(deleting) => service(env)?.delete_with_workspace(&deleting, &protected_workspaces(env)).await?,
2739 Err(_) => worker::console_error!("workspace.deleting {} could not be read", event.id),
2740 }
2741 return Ok(());
2742 }
2743 if event.kind == "workspace.restored" {
2744 match serde_json::from_value::<WorkspaceRestored>(event.data.clone()) {
2745 Ok(restored) => service(env)?.restore_with_workspace(&restored).await?,
2746 Err(_) => worker::console_error!("workspace.restored {} could not be read", event.id),
2747 }
2748 return Ok(());
2749 }
2750 if event.kind == "workspace.deleted" {
2751 match serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) {
2752 Ok(deleted) => service(env)?.purge_workspace(&deleted, &protected_workspaces(env)).await?,
2753 Err(_) => worker::console_error!("workspace.deleted {} could not be read", event.id),
2754 }
2755 return Ok(());
2756 }
2757 // An account deleted or restored: kept contributors that name it
2758 // (or ghost, for a restore) are worked out again, so it shows as
2759 // ghost for its 30 days, and as itself again if restored.
2760 if let Some(needle) = stats::shown_differently(&event.kind, &event.data) {
2761 let db = env.d1("DB")?;
2762 if let Err(error) = stats::rework_naming(&db, &needle).await {
2763 worker::console_error!("{} {}: contributors not marked to be counted again: {error}", event.kind, event.id);
2764 }
2765 return Ok(());
2766 }
2767 if event.kind != "workspace.renamed" {
2768 return Ok(());
2769 }
2770 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
2771 worker::console_error!("workspace.renamed {} could not be read", event.id);
2772 return Ok(());
2773 };
2774 let names: HashMap<String, String> = g1t_kit::call(
2775 identity,
2776 "usernames",
2777 &g1t_contracts::identity::UsernamesArgs {
2778 ids: vec![renamed.workspace_id.clone()],
2779 },
2780 )
2781 .await?;
2782 let current = names
2783 .get(&renamed.workspace_id)
2784 .cloned()
2785 .unwrap_or_else(|| renamed.to.clone());
2786 let left = registry
2787 .rename_namespace(&renamed.stale_slugs(&current), &current)
2788 .await?;
2789 if left > 0 {
2790 worker::console_error!(
2791 "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
2792 renamed.from,
2793 renamed.to
2794 );
2795 }
2796 Ok(())
2797}
2798
2799/// The workspaces whose repositories never go with a deletion, whatever is
2800/// published: `PROTECTED_WORKSPACES` if set here, and Flagon's always.
2801fn protected_workspaces(env: &Env) -> Vec<String> {
2802 let configured = env.var("PROTECTED_WORKSPACES").ok().map(|v| v.to_string());
2803 g1t_contracts::identity::protected_names(configured.as_deref())
2804}
2805
2806/// What a token's own limits say about git on the repository `repo`
2807/// (`owner/name`), before anyone's role is asked: why it is refused, or
2808/// `None`. `source` is the repository a pull request's working copy at
2809/// `repo` belongs to, which a workflow job's token reaches too. `public`
2810/// is whether anyone may read it, and `exists` whether there is one.
2811///
2812/// A job's token and a deploy key reach their own repository only. Its
2813/// scopes decide the rest: `code:read` to read a private repository,
2814/// `code:write` to push, which a read-only deploy key never has. A deploy
2815/// key never makes a repository by pushing to an empty address.
2816pub(crate) fn git_token_refusal(
2817 access: &g1t_contracts::scopes::TokenAccess,
2818 repo: &str,
2819 source: Option<&str>,
2820 write: bool,
2821 public: bool,
2822 exists: bool,
2823) -> Option<String> {
2824 if let Some(refused) = g1t_contracts::scopes::decide_repo(access, repo)
2825 && !source.is_some_and(|source| access.reaches(source))
2826 {
2827 return Some(refused.reason.unwrap_or_default());
2828 }
2829 let decision = g1t_contracts::scopes::decide_git(access, write, public);
2830 if !decision.allowed {
2831 return Some(decision.reason.unwrap_or_default());
2832 }
2833 if access.deploy_key.is_some() && !exists {
2834 return Some(format!("This deploy key is for {repo}, which is not there any more."));
2835 }
2836 None
2837}
2838
2839/// The repository a push to a path that does not exist yet creates: private,
2840/// so nothing pushed by mistake is published. An owner makes it public on
2841/// purpose (`POST /repos/{owner}/{repo}/visibility`).
2842fn push_to_create(owner: &User, path: &RepoPath) -> CreateArgs {
2843 CreateArgs {
2844 owner: owner.clone(),
2845 namespace: path.namespace.clone(),
2846 name: path.name.clone(),
2847 description: None,
2848 is_private: true,
2849 import_url: None,
2850 import_token: None,
2851 }
2852}
2853
2854/// Whether `delete_branch` may remove `branch`: one of g1t's own
2855/// (`g1t-…`), or one whose tip the caller names, such as a dependency
2856/// update's branch after its pull request closed.
2857fn deletable_branch(branch: &str, head: Option<&str>) -> bool {
2858 !branch.is_empty() && (branch.starts_with(G1T_BRANCH_PREFIX) || head.is_some_and(|head| !head.is_empty()))
2859}
2860
2861#[cfg(test)]
2862mod delete_branch_tests {
2863 use super::deletable_branch;
2864
2865 #[test]
2866 fn only_g1t_branches_or_a_named_tip_are_deleted() {
2867 assert!(deletable_branch("g1t-queue-12", None));
2868 assert!(!deletable_branch("g1t/security/sharp-0.35.5", None));
2869 assert!(deletable_branch("g1t/security/sharp-0.35.5", Some("abc123")));
2870 assert!(!deletable_branch("feature", Some("")));
2871 assert!(!deletable_branch("", Some("abc123")));
2872 }
2873}
2874
2875#[cfg(test)]
2876mod announced_refs_tests {
2877 use super::*;
2878
2879 fn pushed(git_ref: String) -> git_http::Pushed {
2880 git_http::Pushed { git_ref, before: None, after: "abc".into() }
2881 }
2882
2883 fn refs(announced: Vec<&git_http::Pushed>) -> Vec<&str> {
2884 announced.into_iter().map(|pushed| pushed.git_ref.as_str()).collect()
2885 }
2886
2887 #[test]
2888 fn a_few_tags_are_announced_and_many_are_not() {
2889 let few: Vec<_> = (1..=3).map(|n| pushed(format!("refs/tags/v{n}"))).chain([pushed("refs/heads/main".into())]).collect();
2890 assert_eq!(refs(announced_refs(&few, "main")), ["refs/heads/main", "refs/tags/v1", "refs/tags/v2", "refs/tags/v3"]);
2891 let many: Vec<_> = (1..=10_000).map(|n| pushed(format!("refs/tags/v{n}"))).chain([pushed("refs/heads/main".into())]).collect();
2892 assert_eq!(refs(announced_refs(&many, "main")), ["refs/heads/main"]);
2893 }
2894
2895 #[test]
2896 fn past_the_branch_cap_only_the_default_branch_is_announced() {
2897 let at_cap: Vec<_> = (0..MAX_PUSH_BRANCH_EVENTS).map(|n| pushed(format!("refs/heads/b{n}"))).collect();
2898 assert_eq!(announced_refs(&at_cap, "main").len(), MAX_PUSH_BRANCH_EVENTS);
2899 let over: Vec<_> = (0..=MAX_PUSH_BRANCH_EVENTS)
2900 .map(|n| pushed(format!("refs/heads/b{n}")))
2901 .chain([pushed("refs/heads/main".into())])
2902 .collect();
2903 assert_eq!(refs(announced_refs(&over, "main")), ["refs/heads/main"]);
2904 assert!(announced_refs(&over, "trunk").is_empty());
2905 }
2906}
2907
2908#[cfg(test)]
2909mod push_to_create_tests {
2910 use super::*;
2911
2912 #[test]
2913 fn a_pushed_repository_starts_private() {
2914 let owner: User = serde_json::from_value(serde_json::json!({ "id": "usr_1", "username": "ada" })).unwrap();
2915 let args = push_to_create(&owner, &RepoPath { namespace: "acme".into(), name: "site".into() });
2916 assert!(args.is_private);
2917 assert_eq!((args.namespace.as_str(), args.name.as_str()), ("acme", "site"));
2918 }
2919}
2920
2921#[cfg(test)]
2922mod deploy_key_git_tests {
2923 use super::*;
2924 use g1t_contracts::deploy_keys;
2925
2926 fn key(read_only: bool) -> User {
2927 deploy_keys::principal("wsp_acme", "acme", deploy_keys::access("dk_1", "CI", "acme/rocket", read_only))
2928 }
2929
2930 fn rocket(private: bool) -> Repo {
2931 serde_json::from_value(serde_json::json!({
2932 "id": "rep_rocket",
2933 "namespace": "acme",
2934 "name": "rocket",
2935 "description": null,
2936 "isPrivate": private,
2937 "ownerId": "usr_owner",
2938 "defaultBranch": "main",
2939 "forkOf": null,
2940 "protected": false,
2941 "createdAt": "",
2942 }))
2943 .unwrap()
2944 }
2945
2946 fn refusal(user: &User, repo: &str, write: bool, exists: bool) -> Option<String> {
2947 git_token_refusal(user.token.as_deref().unwrap(), repo, None, write, false, exists)
2948 }
2949
2950 #[test]
2951 fn a_read_only_deploy_key_clones_its_repository_and_never_pushes() {
2952 let user = key(true);
2953 assert_eq!(refusal(&user, "acme/rocket", false, true), None);
2954 assert!(refusal(&user, "acme/rocket", true, true).unwrap().contains("read-only"));
2955 // Its role is a workspace token's: it reads a private repository.
2956 assert!(registry::can_read(&rocket(true), &Some(user)));
2957 }
2958
2959 #[test]
2960 fn a_deploy_key_with_write_access_pushes_to_its_repository() {
2961 let user = key(false);
2962 assert_eq!(refusal(&user, "acme/rocket", true, true), None);
2963 assert!(registry::can_write(&rocket(true), &Some(user)));
2964 }
2965
2966 #[test]
2967 fn a_deploy_key_reaches_no_other_repository() {
2968 let user = key(false);
2969 for other in ["acme/booster", "other/rocket"] {
2970 for write in [false, true] {
2971 let why = refusal(&user, other, write, true).expect(other);
2972 assert!(why.contains("deploy key is for acme/rocket"), "{why}");
2973 }
2974 }
2975 // Not even a pull request's working copy of another repository.
2976 let token = user.token.as_deref().unwrap();
2977 assert!(git_token_refusal(token, "pulls/pr_1", Some("acme/booster"), false, false, true).is_some());
2978 }
2979
2980 #[test]
2981 fn a_deploy_key_never_creates_a_repository() {
2982 let why = refusal(&key(false), "acme/rocket", true, false).unwrap();
2983 assert!(why.contains("not there"), "{why}");
2984 }
2985}