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