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