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