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