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