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