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