g1t/services/repos/src/lib.rs

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