g1t/services/repos/src/lib.rs

547 lines19,671 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Rust repos service with shipping; pull requests kept in the model1//! 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 git_http;
9mod land;
10mod registry;
11mod store;
12
13use g1t_contracts::events::{GitPush, NewEvent, RepoCreated, RepoForked};
14use g1t_contracts::repos::*;
15use g1t_contracts::{
16 FailureCode, Outcome, User, Viewer, is_valid_namespace, is_valid_repo_name, new_id,
17};
18use g1t_kit::{args, js, now_ms, reply, rpc_method};
19use std::collections::{HashMap, HashSet};
20
21use serde::Serialize;
22use worker::wasm_bindgen::JsValue;
23use worker::{Context, Env, Request, Response, Result, event};
24
25use registry::{Registry, can_read, can_write, store_key};
26use store::{ArtifactsStore, GitRepo, GitStore, Scope};
27
28/// Namespace that holds every attempt's fork: `attempts/<attempt id>`.
29const ATTEMPTS_NAMESPACE: &str = "attempts";
30const MAX_TEXT_BYTES: usize = 512 * 1024;
31/// How far back an attempt may have forked and still be landed.
32const MAX_ANCESTRY: u32 = 1000;
33const SOURCE: &str = "repos";
34const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
35
36fn not_found<T>() -> Outcome<T> {
37 Outcome::fail(FailureCode::NotFound, "Repository not found.")
38}
39
40/// Decoded text, or `None` when the file is too large or looks binary.
41fn text_of(bytes: Vec<u8>) -> Option<String> {
42 if bytes.len() > MAX_TEXT_BYTES || bytes.contains(&0) {
43 return None;
44 }
45 Some(String::from_utf8_lossy(&bytes).into_owned())
46}
47
48fn is_readme(name: &str) -> bool {
49 matches!(
50 name.to_lowercase().as_str(),
51 "readme" | "readme.md" | "readme.markdown" | "readme.txt"
52 )
53}
54
55/// Whether `ancestor` is reachable from the newest commit in `history`.
56///
57/// `history` is the first-parent chain, which is all the store lists; a fork
58/// that merged the target branch in has the target's head on a second
59/// parent, so the walk follows every parent.
60async fn descends_from<R: GitRepo>(repo: &R, history: &[Commit], ancestor: &str) -> Result<bool> {
61 let known: HashMap<&str, &[String]> = history
62 .iter()
63 .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
64 .collect();
65 let mut seen = HashSet::new();
66 let mut queue: Vec<String> = history
67 .first()
68 .map(|c| c.hash.clone())
69 .into_iter()
70 .collect();
71 while let Some(hash) = queue.pop() {
72 if hash == ancestor {
73 return Ok(true);
74 }
75 if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
76 continue;
77 }
78 match known.get(hash.as_str()) {
79 Some(parents) => queue.extend(parents.iter().cloned()),
80 None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
81 }
82 }
83 Ok(false)
84}
85
86struct Repos<S: GitStore> {
87 registry: Registry,
88 store: S,
89 /// The events service, an RPC stub.
90 events: JsValue,
91}
92
93impl<S: GitStore> Repos<S> {
94 async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> {
95 js::call(&self.events, "publish", &[js::to_js(&[event])?]).await?;
96 Ok(())
97 }
98
99 /// Resolves a repo the viewer may read; private repos look missing.
100 async fn readable(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
101 Ok(self
102 .registry
103 .by_path(path)
104 .await?
105 .filter(|repo| can_read(repo, viewer)))
106 }
107
108 async fn get(&self, a: GetArgs) -> Result<Outcome<Repo>> {
109 Ok(self
110 .readable(&a.path, &a.viewer)
111 .await?
112 .map_or_else(not_found, Outcome::Ok))
113 }
114
115 async fn get_by_id(&self, a: GetByIdArgs) -> Result<Outcome<Repo>> {
116 Ok(self
117 .registry
118 .by_id(&a.id)
119 .await?
120 .filter(|repo| can_read(repo, &a.viewer))
121 .map_or_else(not_found, Outcome::Ok))
122 }
123
124 async fn create(&self, a: CreateArgs) -> Result<Outcome<Repo>> {
125 if !a.owner.verified {
126 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
127 }
128 let name = a.name.trim().to_lowercase();
129 if !is_valid_repo_name(&name) {
130 return Ok(Outcome::fail(
131 FailureCode::Invalid,
132 "Use letters, digits, dots, hyphens and underscores only.",
133 ));
134 }
135 if !is_valid_namespace(&a.owner.username) {
136 return Ok(Outcome::fail(
137 FailureCode::Invalid,
138 "This account cannot own repositories.",
139 ));
140 }
141 let path = RepoPath {
142 namespace: a.owner.username.clone(),
143 name,
144 };
145 if self.registry.by_path(&path).await?.is_some() {
146 return Ok(Outcome::fail(
147 FailureCode::Conflict,
148 "You already have a repository with that name.",
149 ));
150 }
151 let now = now_ms();
152 let repo = Repo {
153 id: new_id("rep", now),
154 namespace: path.namespace,
155 name: path.name,
156 description: a
157 .description
158 .map(|text| text.trim().to_owned())
159 .filter(|text| !text.is_empty()),
160 is_private: a.is_private,
161 owner_id: a.owner.id.clone(),
162 default_branch: "main".to_owned(),
163 fork_of: None,
164 created_at: now,
165 };
166 self.store
167 .create(
168 &store_key(&repo),
169 repo.description.as_deref(),
170 &repo.default_branch,
171 )
172 .await?;
173 self.registry.insert(&repo).await?;
174 self.publish(NewEvent {
175 kind: "repo.created",
176 source: SOURCE,
177 repo_id: Some(repo.id.clone()),
178 actor: Some(a.owner.id),
179 data: RepoCreated {
180 repo_id: repo.id.clone(),
181 namespace: repo.namespace.clone(),
182 name: repo.name.clone(),
183 is_private: repo.is_private,
184 },
185 })
186 .await?;
187 Ok(Outcome::Ok(repo))
188 }
189
190 async fn tree(&self, a: TreeArgs) -> Result<Outcome<TreeView>> {
191 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
192 return Ok(not_found());
193 };
194 let git = self.store.open(&store_key(&repo)).await?;
195 let git_ref = a
196 .git_ref
197 .clone()
198 .unwrap_or_else(|| repo.default_branch.clone());
199
200 let Some(head) = git.log(&git_ref, 1).await?.into_iter().next() else {
201 // An unknown ref is an error; a repo with no commits is just empty.
202 if a.git_ref.is_some() {
203 return Ok(Outcome::fail(
204 FailureCode::NotFound,
205 "No such branch, tag or commit.",
206 ));
207 }
208 return Ok(Outcome::Ok(TreeView {
209 repo,
210 git_ref,
211 path: a.tree_path,
212 head: None,
213 entries: Vec::new(),
214 readme: None,
215 }));
216 };
217
218 let no_directory = || Outcome::fail(FailureCode::NotFound, "No such directory.");
219 let mut entries = git.read_tree(&head.tree_hash).await?;
220 for segment in a.tree_path.split('/').filter(|segment| !segment.is_empty()) {
221 let next = entries.as_ref().and_then(|entries| {
222 entries
223 .iter()
224 .find(|entry| entry.name == segment && entry.kind == EntryKind::Tree)
225 });
226 let Some(next) = next else {
227 return Ok(no_directory());
228 };
229 entries = git.read_tree(&next.hash).await?;
230 }
231 let Some(mut entries) = entries else {
232 return Ok(no_directory());
233 };
234 // Directories first, then by name.
235 entries.sort_by(|a, b| {
236 (b.kind == EntryKind::Tree)
237 .cmp(&(a.kind == EntryKind::Tree))
238 .then_with(|| a.name.cmp(&b.name))
239 });
240
241 let readme_entry = entries
242 .iter()
243 .find(|entry| entry.kind == EntryKind::Blob && is_readme(&entry.name));
244 let readme = match readme_entry {
245 Some(entry) => git.read_blob(&entry.hash).await?.map(|bytes| Readme {
246 name: entry.name.clone(),
247 text: text_of(bytes),
248 }),
249 None => None,
250 };
251 Ok(Outcome::Ok(TreeView {
252 repo,
253 git_ref,
254 path: a.tree_path,
255 head: Some(head),
256 entries,
257 readme,
258 }))
259 }
260
261 async fn blob(&self, a: BlobArgs) -> Result<Outcome<BlobView>> {
262 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
263 return Ok(not_found());
264 };
265 let bytes = if a.file_path.is_empty() {
266 None
267 } else {
268 let git = self.store.open(&store_key(&repo)).await?;
269 git.read_file(&a.git_ref, &a.file_path).await?
270 };
271 let Some(bytes) = bytes else {
272 return Ok(Outcome::fail(FailureCode::NotFound, "No such file."));
273 };
274 Ok(Outcome::Ok(BlobView {
275 repo,
276 git_ref: a.git_ref,
277 path: a.file_path,
278 size: bytes.len() as u64,
279 text: text_of(bytes),
280 }))
281 }
282
283 async fn log(&self, a: LogArgs) -> Result<Outcome<Vec<Commit>>> {
284 let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
285 return Ok(not_found());
286 };
287 let git = self.store.open(&store_key(&repo)).await?;
288 let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
289 Ok(Outcome::Ok(git.log(&git_ref, a.limit).await?))
290 }
291
292 async fn fork_for_attempt(&self, a: ForkArgs) -> Result<Outcome<Repo>> {
293 let viewer = Some(a.actor.clone());
294 let Some(source) = self
295 .registry
296 .by_id(&a.source_id)
297 .await?
298 .filter(|repo| can_read(repo, &viewer))
299 else {
300 return Ok(not_found());
301 };
302 let now = now_ms();
303 let fork = Repo {
304 id: new_id("rep", now),
305 namespace: ATTEMPTS_NAMESPACE.to_owned(),
306 name: a.attempt_id.clone(),
307 description: None,
308 // A fork is exactly as visible as the repo it came from.
309 is_private: source.is_private,
310 owner_id: a.actor.id.clone(),
311 default_branch: source.default_branch.clone(),
312 fork_of: Some(source.id.clone()),
313 created_at: now,
314 };
315 self.store
316 .open(&store_key(&source))
317 .await?
318 .fork(&store_key(&fork))
319 .await?;
320 self.registry.insert(&fork).await?;
321 self.publish(NewEvent {
322 kind: "repo.forked",
323 source: SOURCE,
324 repo_id: Some(source.id.clone()),
325 actor: Some(a.actor.id),
326 data: RepoForked {
327 repo_id: fork.id.clone(),
328 source_repo_id: source.id,
329 attempt_id: a.attempt_id,
330 },
331 })
332 .await?;
333 Ok(Outcome::Ok(fork))
334 }
335
336 async fn git_access(&self, a: GitAccessArgs) -> Result<Outcome<GitAccess>> {
337 let write = a.service == GitService::ReceivePack;
338 // Anonymous callers are asked to authenticate whether or not the repo
339 // exists, so private repos cannot be told apart from missing ones.
340 let denied = || match &a.viewer {
341 Some(_) => not_found(),
342 None => Outcome::fail(FailureCode::Unauthenticated, "Authentication required."),
343 };
344 if let (true, Some(user)) = (write, &a.viewer)
345 && !user.verified
346 {
347 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
348 }
349
350 let repo = match self.registry.by_path(&a.path).await? {
351 Some(repo) => {
352 let allowed = if write {
353 can_write(&repo, &a.viewer)
354 } else {
355 can_read(&repo, &a.viewer)
356 };
357 if !allowed {
358 return Ok(denied());
359 }
360 repo
361 }
362 None => {
363 // Push to create, in the pusher's own namespace only.
364 let owner = a
365 .viewer
366 .as_ref()
367 .filter(|user| write && user.username == a.path.namespace.to_lowercase());
368 let Some(owner) = owner else {
369 return Ok(denied());
370 };
371 let created = self
372 .create(CreateArgs {
373 owner: owner.clone(),
374 name: a.path.name.clone(),
375 description: None,
376 is_private: false,
377 })
378 .await?;
379 match created {
380 Outcome::Ok(repo) => repo,
381 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
382 }
383 }
384 };
385 let git = self.store.open(&store_key(&repo)).await?;
386 let scope = if write { Scope::Write } else { Scope::Read };
387 Ok(Outcome::Ok(git.access(scope).await?))
388 }
389
390 async fn land(&self, a: LandArgs) -> Result<Outcome<Landed>> {
391 let actor: Viewer = Some(a.actor.clone());
392 let Some(fork) = self.registry.by_id(&a.fork_id).await? else {
393 return Ok(not_found());
394 };
395 let target = match &fork.fork_of {
396 Some(id) => self.registry.by_id(id).await?,
397 None => None,
398 };
399 let Some(target) = target.filter(|repo| can_read(repo, &actor)) else {
400 return Ok(not_found());
401 };
402 if !can_write(&target, &actor) {
403 return Ok(Outcome::fail(
404 FailureCode::Forbidden,
405 "Only the repository's owner can land an attempt.",
406 ));
407 }
408 if !a.actor.verified {
409 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
410 }
411
412 let branch = &target.default_branch;
413 let fork_git = self.store.open(&store_key(&fork)).await?;
414 let target_git = self.store.open(&store_key(&target)).await?;
415 let history = fork_git.log(branch, MAX_ANCESTRY).await?;
416 let Some(new) = history.first().map(|commit| commit.hash.clone()) else {
417 return Ok(Outcome::fail(
418 FailureCode::Conflict,
419 "This attempt has no commits to land.",
420 ));
421 };
422 let old = target_git
423 .log(branch, 1)
424 .await?
425 .into_iter()
426 .next()
427 .map(|commit| commit.hash);
428
429 if old.as_deref() == Some(new.as_str()) {
430 return Ok(Outcome::Ok(Landed { commit: new }));
431 }
432 // Moving the branch to a commit that does not descend from its
433 // current head would discard whatever landed in between.
434 if let Some(old) = &old
435 && !descends_from(&fork_git, &history, old).await?
436 {
437 return Ok(Outcome::fail(
438 FailureCode::Conflict,
439 format!(
440 "{branch} has moved since this attempt started. Pull {branch} into the attempt's fork, push, and land again."
441 ),
442 ));
443 }
444
445 let source_access = fork_git.access(Scope::Read).await?;
446 let target_access = target_git.access(Scope::Write).await?;
447 let pushed =
448 land::fast_forward(&source_access, &target_access, branch, old.as_deref(), &new)
449 .await?;
450 if let Err(reason) = pushed {
451 // Most often another attempt landed between the check and the push.
452 return Ok(Outcome::fail(
453 FailureCode::Conflict,
454 format!("{branch} could not be updated: {reason}"),
455 ));
456 }
457 self.publish_push(&target, &new, Some(a.actor.id)).await?;
458 Ok(Outcome::Ok(Landed { commit: new }))
459 }
460
461 async fn publish_push(&self, repo: &Repo, after: &str, actor: Option<String>) -> Result<()> {
462 self.publish(NewEvent {
463 kind: "git.push",
464 source: SOURCE,
465 repo_id: Some(repo.id.clone()),
466 actor,
467 data: GitPush {
468 repo_id: repo.id.clone(),
469 git_ref: format!("refs/heads/{}", repo.default_branch),
470 after: after.to_owned(),
471 },
472 })
473 .await
474 }
475
476 /// Git over HTTPS.
477 async fn git_http(&self, request: Request, env: &Env) -> Result<Response> {
478 let Some(git) = git_http::parse(&request.url()?) else {
479 return Response::error("Not found", 404);
480 };
481 let viewer = git_http::viewer(&request, &env.service("IDENTITY")?).await?;
482 let access = self
483 .git_access(GitAccessArgs {
484 path: git.path.clone(),
485 viewer: viewer.clone(),
486 service: git.service,
487 })
488 .await?;
489 let access = match access {
490 Outcome::Ok(access) => access,
491 refused => return git_http::refuse(refused),
492 };
493 let response = git_http::forward(request, &git, &access).await?;
494
495 // Artifacts' own push notifications are per repository, which does
496 // not fit a repo per attempt, so the front end reports pushes itself.
497 let pushed = git.endpoint == "git-receive-pack" && response.status_code() == 200;
498 if pushed && let Some(repo) = self.registry.by_path(&git.path).await? {
499 let head = self
500 .store
501 .open(&store_key(&repo))
502 .await?
503 .log(&repo.default_branch, 1)
504 .await?;
505 if let Some(head) = head.first() {
506 self.publish_push(&repo, &head.hash, viewer.map(|user: User| user.id))
507 .await?;
508 }
509 }
510 Ok(response)
511 }
512}
513
514#[event(fetch)]
515async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
516 let repos = Repos {
517 registry: Registry { db: env.d1("DB")? },
518 store: ArtifactsStore::new(&env)?,
519 events: js::binding(&env, "EVENTS")?,
520 };
521 let Some(method) = rpc_method(&request) else {
522 return repos.git_http(request, &env).await;
523 };
524 let body: serde_json::Value = request.json().await?;
525
526 match method.as_str() {
527 "get" => reply(&repos.get(args(body)?).await?),
528 "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
529 "list" => {
530 let a: ListArgs = args(body)?;
531 reply(
532 &repos
533 .registry
534 .list(&a.viewer, a.query.as_deref(), a.namespace.as_deref())
535 .await?,
536 )
537 }
538 "create" => reply(&repos.create(args(body)?).await?),
539 "tree" => reply(&repos.tree(args(body)?).await?),
540 "blob" => reply(&repos.blob(args(body)?).await?),
541 "log" => reply(&repos.log(args(body)?).await?),
542 "fork_for_attempt" => reply(&repos.fork_for_attempt(args(body)?).await?),
543 "git_access" => reply(&repos.git_access(args(body)?).await?),
544 "land" => reply(&repos.land(args(body)?).await?),
545 _ => Response::error("Unknown method", 404),
546 }
547}