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