g1t/services/repos/src/lib.rs

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