This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-xrpc / src / repos.rs
22 kB 668 lines
1use std::sync::Arc; 2 3use axum::Json; 4use axum::body::Bytes; 5use axum::extract::State; 6use axum::response::{IntoResponse, Response}; 7use http::HeaderMap; 8use serde::{Deserialize, Serialize}; 9 10use knot_acl::{KnotAcl, can_admin_knot, can_create_repo, can_delete_repo}; 11use knot_atproto::{PreparedRepoDid, RecordPresence}; 12use knot_cob::{CobHome, CobStore}; 13use knot_cobs::{ 14 Registration, RegistryChange, Rename, RepoRef, RepoRegistryCob, deregister_repo, register_repo, 15}; 16use knot_git::{GitError, Layout, Repo}; 17use knot_index::Resolved; 18use knot_runtime::{Clock, HttpTransport, Signer}; 19use knot_types::{ 20 AccountDid, ActorId, BranchName, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, UnixSeconds, 21}; 22 23use crate::body::SourceUrl; 24use crate::error::XrpcError; 25use crate::reservations::ReserveDecision; 26use crate::{XrpcState, decode, ok_empty, run_blocking}; 27 28pub(crate) const CREATE_ROUTE: &str = "/xrpc/sh.tangled.repo.create"; 29pub(crate) const DELETE_ROUTE: &str = "/xrpc/sh.tangled.repo.delete"; 30pub(crate) const RENAME_ROUTE: &str = "/xrpc/sh.tangled.repo.rename"; 31pub(crate) const RESERVE_ROUTE: &str = "/xrpc/sh.tangled.repo.reserveKey"; 32 33#[derive(Deserialize)] 34struct CreateInput { 35 rkey: RepoRkey, 36 name: RepoName, 37 #[serde(rename = "defaultBranch")] 38 default_branch: Option<BranchName>, 39 #[serde(default, deserialize_with = "crate::body::optional_source_url")] 40 source: Option<SourceUrl>, 41 #[serde(rename = "repoDid")] 42 repo_did: Option<RepoDid>, 43} 44 45#[derive(Serialize)] 46struct CreateOutput { 47 #[serde(rename = "repoDid")] 48 repo_did: RepoDid, 49 key: ActorId, 50 #[serde(rename = "lfsMissing", skip_serializing_if = "Vec::is_empty")] 51 lfs_missing: Vec<knot_lfs::LfsOid>, 52} 53 54#[derive(Deserialize)] 55struct DeleteInput { 56 did: OwnerDid, 57 rkey: RepoRkey, 58 #[serde(rename = "name")] 59 _name: RepoName, 60 #[serde(default)] 61 force: bool, 62} 63 64#[derive(Deserialize)] 65struct RenameInput { 66 repo: RepoDid, 67 rkey: RepoRkey, 68 name: RepoName, 69} 70 71#[derive(Deserialize)] 72struct ReserveInput { 73 #[serde(rename = "repoDid")] 74 repo_did: RepoDid, 75} 76 77#[derive(Serialize)] 78struct ReserveOutput { 79 #[serde(rename = "repoDid")] 80 repo_did: RepoDid, 81 key: ActorId, 82} 83 84pub(crate) async fn reserve_key<H: HttpTransport, C: Clock>( 85 State(state): State<Arc<XrpcState<H, C>>>, 86 headers: HeaderMap, 87 method: crate::Method, 88 body: Bytes, 89) -> Result<Response, XrpcError> { 90 let actor = state.authenticate(&headers, &method).await?; 91 let acl = KnotAcl::new(&state.admins, state.admission, &state.index); 92 if !can_create_repo(&acl, &actor).is_allowed() { 93 return Err(XrpcError::forbidden( 94 "only knot admin or member may reserve a repository key", 95 )); 96 } 97 98 let ReserveInput { repo_did } = decode(&body)?; 99 if !repo_did.as_str().starts_with("did:web:") { 100 return Err(XrpcError::invalid_request( 101 "only did:web repo identity needs a reserved key. Omit repoDid on create to mint a did:plc.", 102 )); 103 } 104 if repo_did.as_str() == state.knot_did.as_str() { 105 return Err(XrpcError::invalid_request( 106 "repoDid mustn't be knot's own identity", 107 )); 108 } 109 match state.index.owner_of(&repo_did) { 110 Resolved::Ready(Some(_)) => { 111 return Err(XrpcError::conflict( 112 "that repo DID is already hosted on this knot", 113 )); 114 } 115 Resolved::Warming => { 116 return Err(XrpcError::warming("registry projection is still warming")); 117 } 118 Resolved::Ready(None) => {} 119 } 120 121 let now = state.now(); 122 state.reservations.prune(now); 123 124 match state.reservations.try_reserve(&repo_did, &actor, now) { 125 ReserveDecision::HeldByOther => { 126 return Err(XrpcError::conflict( 127 "that repo DID is reserved by another account", 128 )); 129 } 130 ReserveDecision::PerActorFull => { 131 return Err(XrpcError::rate_limited( 132 "you are holding maximum number of reserved repository keys awaiting creation", 133 )); 134 } 135 ReserveDecision::GlobalFull => { 136 return Err(XrpcError::rate_limited( 137 "knot is holding maximum number of reserved repository keys awaiting creation", 138 )); 139 } 140 ReserveDecision::Fresh | ReserveDecision::Renewed => {} 141 } 142 143 let public = match state.secrets.public_key(&state.knot_did) { 144 Ok(public) => public, 145 Err(error) => { 146 state.reservations.release(&repo_did); 147 return Err(error.into()); 148 } 149 }; 150 let key = ActorId::from_secp256k1(public.as_bytes()); 151 152 Ok((http::StatusCode::OK, Json(ReserveOutput { repo_did, key })).into_response()) 153} 154 155pub(crate) async fn create_repo<H: HttpTransport, C: Clock>( 156 State(state): State<Arc<XrpcState<H, C>>>, 157 headers: HeaderMap, 158 method: crate::Method, 159 body: Bytes, 160) -> Result<Response, XrpcError> { 161 let actor = state.authenticate(&headers, &method).await?; 162 let acl = KnotAcl::new(&state.admins, state.admission, &state.index); 163 if !can_create_repo(&acl, &actor).is_allowed() { 164 return Err(XrpcError::forbidden( 165 "only knot admin or member may create repositories", 166 )); 167 } 168 169 let input: CreateInput = decode(&body)?; 170 let source = input 171 .source 172 .as_ref() 173 .map(|source| { 174 crate::forks::resolve_upstream(&state, source).map(|upstream| (source, upstream)) 175 }) 176 .transpose()?; 177 let head = input.default_branch.map(|branch| branch.head_ref()); 178 179 let owner = OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID"); 180 match state.index.resolve_repo(&owner, &input.rkey) { 181 Resolved::Ready(Some(existing)) => match state.index.rkey_of(&existing) { 182 Resolved::Ready(Some(canonical)) if canonical == input.rkey => { 183 return Err(XrpcError::conflict( 184 "repository with that record key already exists for this owner", 185 )); 186 } 187 Resolved::Warming => { 188 return Err(XrpcError::warming("registry projection is still warming")); 189 } 190 _ => {} 191 }, 192 Resolved::Warming => { 193 return Err(XrpcError::warming("registry projection is still warming")); 194 } 195 Resolved::Ready(None) => {} 196 } 197 198 let (repo_did, provisioning) = 199 provision_repo_did(&state, &actor, &input.rkey, input.repo_did).await?; 200 let now = state.now(); 201 let knot_signer = state.secrets.signer(&state.knot_did)?; 202 let key = ActorId::from_secp256k1(knot_signer.public_key().as_bytes()); 203 204 let registration = Registration { 205 owner, 206 rkey: input.rkey, 207 name: input.name, 208 repo: repo_did.clone(), 209 created_at: now, 210 }; 211 212 let lfs_store = state.lfs.as_ref().map(|web| Arc::clone(&web.handle.store)); 213 let layout = state.layout.clone(); 214 let placed = repo_did.clone(); 215 let lfs = lfs_store.clone(); 216 let provisioning = run_blocking(move || { 217 let repo = layout.create(&placed).map_err(|error| match error { 218 GitError::AlreadyExists(_) => XrpcError::conflict("repository already exists on disk"), 219 GitError::ReservedDid(_) => { 220 XrpcError::invalid_request("repoDid mustn't be knot's own identity") 221 } 222 other => XrpcError::internal(other.to_string()), 223 })?; 224 let outcome = stage_repo(&repo, head.as_ref()); 225 if outcome.is_err() { 226 rollback_local(&layout, lfs.as_deref(), &placed); 227 } 228 outcome.map(|()| provisioning) 229 }) 230 .await?; 231 232 let lfs_missing = match &source { 233 Some((source_url, upstream)) => { 234 match populate(&state, &repo_did, source_url, upstream).await { 235 Ok(missing) => missing, 236 Err(error) => { 237 let layout = state.layout.clone(); 238 let placed = repo_did.clone(); 239 let lfs = lfs_store.clone(); 240 let _ = run_blocking(move || { 241 rollback_local(&layout, lfs.as_deref(), &placed); 242 Ok(()) 243 }) 244 .await; 245 return Err(error); 246 } 247 } 248 } 249 None => Vec::new(), 250 }; 251 252 let submission = match &provisioning { 253 Provisioned::Minted(prepared) => state 254 .atproto 255 .submit_plc_operation(prepared) 256 .await 257 .map(|_| ()), 258 Provisioned::Reserved => Ok(()), 259 }; 260 if let Err(error) = submission { 261 let layout = state.layout.clone(); 262 let placed = repo_did.clone(); 263 let lfs = lfs_store.clone(); 264 let _ = run_blocking(move || { 265 rollback_local(&layout, lfs.as_deref(), &placed); 266 Ok(()) 267 }) 268 .await; 269 return Err(error.into()); 270 } 271 272 let reserved = matches!(provisioning, Provisioned::Reserved); 273 let layout = state.layout.clone(); 274 let cob_locks = Arc::clone(&state.cob_locks); 275 let meta_path = state.meta_path.clone(); 276 let placed = repo_did.clone(); 277 let lfs = lfs_store.clone(); 278 let home = CobHome::from(&state.knot_did); 279 run_blocking(move || { 280 let outcome = { 281 let _guard = cob_locks.meta(); 282 Repo::open(&meta_path) 283 .map_err(XrpcError::from) 284 .and_then(|meta| { 285 register( 286 &CobStore::new(&meta), 287 &home, 288 registration, 289 &knot_signer, 290 now, 291 ) 292 }) 293 }; 294 if outcome.is_err() { 295 rollback_local(&layout, lfs.as_deref(), &placed); 296 if !reserved { 297 tracing::error!( 298 repo = %placed, 299 "registration failed after did:plc submitted to PLC directory" 300 ); 301 } 302 } 303 outcome 304 }) 305 .await?; 306 307 if reserved { 308 state.reservations.release(&repo_did); 309 } 310 311 let index = Arc::clone(&state.index); 312 let refreshed = repo_did.clone(); 313 run_blocking(move || { 314 index.refresh_registry().map_err(|error| { 315 XrpcError::internal(format!( 316 "repo {refreshed} was created and registered but registry projection refresh failed: {error}" 317 )) 318 }) 319 }) 320 .await?; 321 322 Ok(( 323 http::StatusCode::OK, 324 Json(CreateOutput { 325 repo_did, 326 key, 327 lfs_missing, 328 }), 329 ) 330 .into_response()) 331} 332 333async fn populate<H: HttpTransport, C: Clock>( 334 state: &Arc<XrpcState<H, C>>, 335 repo_did: &RepoDid, 336 source: &SourceUrl, 337 upstream: &crate::forks::Upstream, 338) -> Result<Vec<knot_lfs::LfsOid>, XrpcError> { 339 let prefixes = vec![ 340 "HEAD".to_string(), 341 "refs/heads/".to_string(), 342 "refs/tags/".to_string(), 343 ]; 344 let refs = crate::forks::upstream_refs(state, upstream, prefixes).await?; 345 let tips = refs.tips(); 346 let pack = crate::forks::upstream_pack( 347 state, 348 upstream, 349 knot_pack::WantOids::new(refs.tips()), 350 knot_pack::HaveOids::default(), 351 ) 352 .await?; 353 let layout = state.layout.clone(); 354 let placed = repo_did.clone(); 355 let origin = source.clone(); 356 run_blocking(move || { 357 let repo = layout.open(&placed)?; 358 crate::forks::populate_fork(&repo, &refs, &pack, &origin) 359 }) 360 .await?; 361 crate::lfs::mirror_fork_objects( 362 Arc::clone(state), 363 upstream.clone(), 364 repo_did.clone(), 365 knot_pack::WantOids::new(tips), 366 knot_pack::HaveOids::default(), 367 ) 368 .await 369} 370 371fn rollback_local(layout: &Layout, lfs: Option<&knot_lfs::DiskStore>, repo_did: &RepoDid) { 372 if let Some(store) = lfs 373 && let Err(error) = store.remove_repo(repo_did) 374 { 375 tracing::error!(repo = %repo_did, %error, "couldn't roll back lfs prefix"); 376 } 377 if let Err(error) = layout.remove(repo_did) { 378 tracing::error!(repo = %repo_did, %error, "couldn't roll back on-disk repo :3"); 379 } 380} 381 382enum Provisioned { 383 Minted(PreparedRepoDid), 384 Reserved, 385} 386 387async fn provision_repo_did<H: HttpTransport, C: Clock>( 388 state: &XrpcState<H, C>, 389 actor: &knot_types::AccountDid, 390 rkey: &RepoRkey, 391 provided: Option<RepoDid>, 392) -> Result<(RepoDid, Provisioned), XrpcError> { 393 match provided { 394 Some(did) if did.as_str().starts_with("did:web:") => { 395 match state.index.owner_of(&did) { 396 Resolved::Ready(Some(_)) => { 397 return Err(XrpcError::conflict( 398 "that repo DID is already hosted on this knot", 399 )); 400 } 401 Resolved::Warming => { 402 return Err(XrpcError::warming("registry projection is still warming")); 403 } 404 Resolved::Ready(None) => {} 405 } 406 if !state.reservations.holder_is(&did, actor, state.now()) { 407 return Err(XrpcError::invalid_request( 408 "reserve this did:web for your own account via sh.tangled.repo.reserveKey before creating it", 409 )); 410 } 411 let knot_public = state.secrets.public_key(&state.knot_did)?; 412 state 413 .atproto 414 .verify_did_web_publishes_key(&did, &knot_public) 415 .await 416 .map_err(XrpcError::from)?; 417 Ok((did, Provisioned::Reserved)) 418 } 419 Some(_) => Err(XrpcError::invalid_request( 420 "repoDid must be did:web hosted on your own domain. Omit it to mint a did:plc.", 421 )), 422 None => { 423 let knot_signer = state.secrets.signer(&state.knot_did)?; 424 let owner = 425 OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID"); 426 let nonce = knot_atproto::MintNonce::mint(&*state.entropy, &owner, rkey); 427 let prepared = 428 knot_atproto::prepare_repo_did(&knot_signer, &state.knot_service_url, &nonce) 429 .map_err(XrpcError::from)?; 430 Ok((prepared.did.clone(), Provisioned::Minted(prepared))) 431 } 432 } 433} 434 435fn stage_repo(repo: &Repo, head: Option<&RefName>) -> Result<(), XrpcError> { 436 if let Some(refname) = head { 437 repo.set_head(refname)?; 438 } 439 Ok(()) 440} 441 442fn register( 443 store: &CobStore, 444 home: &CobHome, 445 registration: Registration, 446 signer: &dyn Signer, 447 now: UnixSeconds, 448) -> Result<(), XrpcError> { 449 match store 450 .list::<RepoRegistryCob>() 451 .map_err(XrpcError::from)? 452 .as_slice() 453 { 454 [] => store 455 .create(home, &RegistryChange::Register(registration), signer, now) 456 .map(|_| ()) 457 .map_err(XrpcError::from), 458 [object] => register_repo(store, home, *object, registration, signer, now) 459 .map(|_| ()) 460 .map_err(XrpcError::from), 461 many => Err(XrpcError::internal(format!( 462 "{} repo registry objects share meta-repo", 463 many.len() 464 ))), 465 } 466} 467 468pub(crate) async fn delete_repo<H: HttpTransport, C: Clock>( 469 State(state): State<Arc<XrpcState<H, C>>>, 470 headers: HeaderMap, 471 method: crate::Method, 472 body: Bytes, 473) -> Result<Response, XrpcError> { 474 let actor = state.authenticate(&headers, &method).await?; 475 let DeleteInput { 476 did, rkey, force, .. 477 } = decode(&body)?; 478 479 let repo_did = match state.index.resolve_repo(&did, &rkey) { 480 Resolved::Ready(Some(repo_did)) => repo_did, 481 Resolved::Ready(None) => return Ok(ok_empty()), 482 Resolved::Warming => { 483 return Err(XrpcError::warming("registry projection is still warming")); 484 } 485 }; 486 487 let acl = KnotAcl::new(&state.admins, state.admission, &state.index); 488 if !can_delete_repo(&acl, &actor, &repo_did).is_allowed() { 489 return Err(XrpcError::forbidden( 490 "only repository owner or a knot admin may delete it", 491 )); 492 } 493 494 if force { 495 if !can_admin_knot(&acl, &actor).is_allowed() { 496 return Err(XrpcError::forbidden( 497 "only knot admin may force a delete past the PDS record check", 498 )); 499 } 500 } else { 501 let owner = AccountDid::from(did.clone()); 502 match state.atproto.repo_record_present(&owner, &rkey).await { 503 Ok(RecordPresence::Present) => { 504 return Err(XrpcError::conflict( 505 "sh.tangled.repo record still exists on the owner's PDS. Remove it there first or force the delete.", 506 )); 507 } 508 Ok(RecordPresence::Absent) => {} 509 Err(error) => { 510 tracing::warn!( 511 repo = %repo_did, 512 %error, 513 "proceeding w/ best-effort delete despite unconfirmed owner PDS record :p" 514 ); 515 } 516 } 517 } 518 519 let now = state.now(); 520 let knot_signer = state.secrets.signer(&state.knot_did)?; 521 let layout = state.layout.clone(); 522 let index = Arc::clone(&state.index); 523 let cob_locks = Arc::clone(&state.cob_locks); 524 let meta_path = state.meta_path.clone(); 525 let target = RepoRef { owner: did, rkey }; 526 let deleted = repo_did.clone(); 527 let home = CobHome::from(&state.knot_did); 528 let lfs_store = state.lfs.as_ref().map(|lfs| Arc::clone(&lfs.handle.store)); 529 run_blocking(move || { 530 let _repo_guard = cob_locks.repo(&deleted); 531 let _meta_guard = cob_locks.meta(); 532 let meta = Repo::open(&meta_path)?; 533 let store = CobStore::new(&meta); 534 deregister(&store, &home, target, &deleted, &knot_signer, now)?; 535 let removal = layout.remove(&deleted); 536 index.refresh_registry().map_err(|error| { 537 XrpcError::internal(format!( 538 "repo {deleted} was deregistered but registry projection refresh failed: {error}" 539 )) 540 })?; 541 if let Some(store) = &lfs_store 542 && let Err(error) = store.remove_repo(&deleted) 543 { 544 tracing::warn!( 545 repo = %deleted, 546 %error, 547 "lfs prefix removal on delete failed, orphan sweep will reclaim" 548 ); 549 } 550 removal.map_err(|error| { 551 XrpcError::internal(format!( 552 "repo {deleted} was deregistered but its on-disk directory couldn't be removed: {error}" 553 )) 554 }) 555 }) 556 .await?; 557 558 Ok(ok_empty()) 559} 560 561fn deregister( 562 store: &CobStore, 563 home: &CobHome, 564 target: RepoRef, 565 expected: &RepoDid, 566 signer: &dyn Signer, 567 now: UnixSeconds, 568) -> Result<(), XrpcError> { 569 match store 570 .list::<RepoRegistryCob>() 571 .map_err(XrpcError::from)? 572 .as_slice() 573 { 574 [] => Ok(()), 575 [object] => deregister_repo(store, home, *object, target, expected.clone(), signer, now) 576 .map(|_| ()) 577 .map_err(XrpcError::from), 578 many => Err(XrpcError::internal(format!( 579 "{} repo registry objects share meta-repo", 580 many.len() 581 ))), 582 } 583} 584 585pub(crate) async fn rename_repo<H: HttpTransport, C: Clock>( 586 State(state): State<Arc<XrpcState<H, C>>>, 587 headers: HeaderMap, 588 method: crate::Method, 589 body: Bytes, 590) -> Result<Response, XrpcError> { 591 let actor = state.authenticate(&headers, &method).await?; 592 let RenameInput { repo, rkey, name } = decode(&body)?; 593 594 let owner = match state.index.owner_of(&repo) { 595 Resolved::Ready(Some(owner)) => owner, 596 Resolved::Ready(None) => { 597 return Err(XrpcError::not_found("no such repository on this knot")); 598 } 599 Resolved::Warming => { 600 return Err(XrpcError::warming("registry projection is still warming")); 601 } 602 }; 603 604 crate::authorize_push( 605 &state, 606 &actor, 607 &repo, 608 "only repository owner or a collaborator may rename it", 609 ) 610 .await?; 611 612 let now = state.now(); 613 let knot_signer = state.secrets.signer(&state.knot_did)?; 614 let index = Arc::clone(&state.index); 615 let cob_locks = Arc::clone(&state.cob_locks); 616 let meta_path = state.meta_path.clone(); 617 let renamed = repo.clone(); 618 let home = CobHome::from(&state.knot_did); 619 run_blocking(move || { 620 { 621 let _guard = cob_locks.meta(); 622 let meta = Repo::open(&meta_path)?; 623 let store = CobStore::new(&meta); 624 match store 625 .list::<RepoRegistryCob>() 626 .map_err(XrpcError::from)? 627 .as_slice() 628 { 629 [] => { 630 return Err(XrpcError::not_found( 631 "no repositories are registered on this knot", 632 )); 633 } 634 [object] => { 635 knot_cobs::rename_repo( 636 &store, 637 &home, 638 *object, 639 Rename { 640 owner, 641 rkey, 642 name, 643 repo: renamed.clone(), 644 }, 645 &knot_signer, 646 now, 647 ) 648 .map(|_| ()) 649 .map_err(XrpcError::from)?; 650 } 651 many => { 652 return Err(XrpcError::internal(format!( 653 "{} repo registry objects share meta-repo", 654 many.len() 655 ))); 656 } 657 } 658 } 659 index.refresh_registry().map_err(|error| { 660 XrpcError::internal(format!( 661 "repo {renamed} was renamed but registry projection refresh failed: {error}" 662 )) 663 }) 664 }) 665 .await?; 666 667 Ok(ok_empty()) 668}