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 663 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 head = input.default_branch.map(|branch| branch.head_ref()); 171 172 let owner = OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID"); 173 match state.index.resolve_repo(&owner, &input.rkey) { 174 Resolved::Ready(Some(existing)) => match state.index.rkey_of(&existing) { 175 Resolved::Ready(Some(canonical)) if canonical == input.rkey => { 176 return Err(XrpcError::conflict( 177 "repository with that record key already exists for this owner", 178 )); 179 } 180 Resolved::Warming => { 181 return Err(XrpcError::warming("registry projection is still warming")); 182 } 183 _ => {} 184 }, 185 Resolved::Warming => { 186 return Err(XrpcError::warming("registry projection is still warming")); 187 } 188 Resolved::Ready(None) => {} 189 } 190 191 let fork = match &input.source { 192 Some(origin) => Some(crate::forks::ForkSource::resolve(&state, origin).await?), 193 None => None, 194 }; 195 let object_format = fork.as_ref().map(|fork| fork.refs.object_format); 196 197 let (repo_did, provisioning) = 198 provision_repo_did(&state, &actor, &input.rkey, input.repo_did).await?; 199 let now = state.now(); 200 let knot_signer = state.secrets.signer(&state.knot_did)?; 201 let key = ActorId::from_secp256k1(knot_signer.public_key().as_bytes()); 202 203 let registration = Registration { 204 owner, 205 rkey: input.rkey, 206 name: input.name, 207 repo: repo_did.clone(), 208 created_at: now, 209 }; 210 211 let lfs_store = state.lfs.as_ref().map(|web| Arc::clone(&web.handle.store)); 212 let layout = state.layout.clone(); 213 let placed = repo_did.clone(); 214 let lfs = lfs_store.clone(); 215 let provisioning = run_blocking(move || { 216 let repo = match object_format { 217 Some(format) => layout.create_with_format(&placed, format), 218 None => layout.create(&placed), 219 } 220 .map_err(|error| match error { 221 GitError::AlreadyExists(_) => XrpcError::conflict("repository already exists on disk"), 222 GitError::ReservedDid(_) => { 223 XrpcError::invalid_request("repoDid mustn't be knot's own identity") 224 } 225 other => XrpcError::internal(other.to_string()), 226 })?; 227 let outcome = stage_repo(&repo, head.as_ref()); 228 if outcome.is_err() { 229 rollback_local(&layout, lfs.as_deref(), &placed); 230 } 231 outcome.map(|()| provisioning) 232 }) 233 .await?; 234 235 let lfs_missing = match &fork { 236 Some(fork) => match populate(&state, &repo_did, fork).await { 237 Ok(missing) => missing, 238 Err(error) => { 239 let layout = state.layout.clone(); 240 let placed = repo_did.clone(); 241 let lfs = lfs_store.clone(); 242 let _ = run_blocking(move || { 243 rollback_local(&layout, lfs.as_deref(), &placed); 244 Ok(()) 245 }) 246 .await; 247 return Err(error); 248 } 249 }, 250 None => Vec::new(), 251 }; 252 253 let submission = match &provisioning { 254 Provisioned::Minted(prepared) => state 255 .atproto 256 .submit_plc_operation(prepared) 257 .await 258 .map(|_| ()), 259 Provisioned::Reserved => Ok(()), 260 }; 261 if let Err(error) = submission { 262 let layout = state.layout.clone(); 263 let placed = repo_did.clone(); 264 let lfs = lfs_store.clone(); 265 let _ = run_blocking(move || { 266 rollback_local(&layout, lfs.as_deref(), &placed); 267 Ok(()) 268 }) 269 .await; 270 return Err(error.into()); 271 } 272 273 let reserved = matches!(provisioning, Provisioned::Reserved); 274 let layout = state.layout.clone(); 275 let cob_locks = Arc::clone(&state.cob_locks); 276 let meta_path = state.meta_path.clone(); 277 let placed = repo_did.clone(); 278 let lfs = lfs_store.clone(); 279 let home = CobHome::from(&state.knot_did); 280 run_blocking(move || { 281 let outcome = { 282 let _guard = cob_locks.meta(); 283 Repo::open(&meta_path) 284 .map_err(XrpcError::from) 285 .and_then(|meta| { 286 register( 287 &CobStore::new(&meta), 288 &home, 289 registration, 290 &knot_signer, 291 now, 292 ) 293 }) 294 }; 295 if outcome.is_err() { 296 rollback_local(&layout, lfs.as_deref(), &placed); 297 if !reserved { 298 tracing::error!( 299 repo = %placed, 300 "registration failed after did:plc submitted to PLC directory" 301 ); 302 } 303 } 304 outcome 305 }) 306 .await?; 307 308 if reserved { 309 state.reservations.release(&repo_did); 310 } 311 312 let index = Arc::clone(&state.index); 313 let refreshed = repo_did.clone(); 314 run_blocking(move || { 315 index.refresh_registry().map_err(|error| { 316 XrpcError::internal(format!( 317 "repo {refreshed} was created and registered but registry projection refresh failed: {error}" 318 )) 319 }) 320 }) 321 .await?; 322 323 Ok(( 324 http::StatusCode::OK, 325 Json(CreateOutput { 326 repo_did, 327 key, 328 lfs_missing, 329 }), 330 ) 331 .into_response()) 332} 333 334async fn populate<H: HttpTransport, C: Clock>( 335 state: &Arc<XrpcState<H, C>>, 336 repo_did: &RepoDid, 337 fork: &crate::forks::ForkSource, 338) -> Result<Vec<knot_lfs::LfsOid>, XrpcError> { 339 let tips = fork.refs.tips(); 340 let pack = crate::forks::upstream_pack( 341 state, 342 &fork.upstream, 343 knot_pack::WantOids::new(tips.clone()), 344 knot_pack::HaveOids::default(), 345 ) 346 .await?; 347 let layout = state.layout.clone(); 348 let placed = repo_did.clone(); 349 let origin = fork.origin.clone(); 350 let refs = fork.refs.clone(); 351 run_blocking(move || { 352 let repo = layout.open(&placed)?; 353 crate::forks::populate_fork(&repo, &refs, &pack, &origin) 354 }) 355 .await?; 356 crate::lfs::mirror_fork_objects( 357 Arc::clone(state), 358 fork.upstream.clone(), 359 repo_did.clone(), 360 knot_pack::WantOids::new(tips), 361 knot_pack::HaveOids::default(), 362 ) 363 .await 364} 365 366fn rollback_local(layout: &Layout, lfs: Option<&knot_lfs::DiskStore>, repo_did: &RepoDid) { 367 if let Some(store) = lfs 368 && let Err(error) = store.remove_repo(repo_did) 369 { 370 tracing::error!(repo = %repo_did, %error, "couldn't roll back lfs prefix"); 371 } 372 if let Err(error) = layout.remove(repo_did) { 373 tracing::error!(repo = %repo_did, %error, "couldn't roll back on-disk repo :3"); 374 } 375} 376 377enum Provisioned { 378 Minted(PreparedRepoDid), 379 Reserved, 380} 381 382async fn provision_repo_did<H: HttpTransport, C: Clock>( 383 state: &XrpcState<H, C>, 384 actor: &knot_types::AccountDid, 385 rkey: &RepoRkey, 386 provided: Option<RepoDid>, 387) -> Result<(RepoDid, Provisioned), XrpcError> { 388 match provided { 389 Some(did) if did.as_str().starts_with("did:web:") => { 390 match state.index.owner_of(&did) { 391 Resolved::Ready(Some(_)) => { 392 return Err(XrpcError::conflict( 393 "that repo DID is already hosted on this knot", 394 )); 395 } 396 Resolved::Warming => { 397 return Err(XrpcError::warming("registry projection is still warming")); 398 } 399 Resolved::Ready(None) => {} 400 } 401 if !state.reservations.holder_is(&did, actor, state.now()) { 402 return Err(XrpcError::invalid_request( 403 "reserve this did:web for your own account via sh.tangled.repo.reserveKey before creating it", 404 )); 405 } 406 let knot_public = state.secrets.public_key(&state.knot_did)?; 407 state 408 .atproto 409 .verify_did_web_publishes_key(&did, &knot_public) 410 .await 411 .map_err(XrpcError::from)?; 412 Ok((did, Provisioned::Reserved)) 413 } 414 Some(_) => Err(XrpcError::invalid_request( 415 "repoDid must be did:web hosted on your own domain. Omit it to mint a did:plc.", 416 )), 417 None => { 418 let knot_signer = state.secrets.signer(&state.knot_did)?; 419 let owner = 420 OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID"); 421 let nonce = knot_atproto::MintNonce::mint(&*state.entropy, &owner, rkey); 422 let prepared = 423 knot_atproto::prepare_repo_did(&knot_signer, &state.knot_service_url, &nonce) 424 .map_err(XrpcError::from)?; 425 Ok((prepared.did.clone(), Provisioned::Minted(prepared))) 426 } 427 } 428} 429 430fn stage_repo(repo: &Repo, head: Option<&RefName>) -> Result<(), XrpcError> { 431 if let Some(refname) = head { 432 repo.set_head(refname)?; 433 } 434 Ok(()) 435} 436 437fn register( 438 store: &CobStore, 439 home: &CobHome, 440 registration: Registration, 441 signer: &dyn Signer, 442 now: UnixSeconds, 443) -> Result<(), XrpcError> { 444 match store 445 .list::<RepoRegistryCob>() 446 .map_err(XrpcError::from)? 447 .as_slice() 448 { 449 [] => store 450 .create(home, &RegistryChange::Register(registration), signer, now) 451 .map(|_| ()) 452 .map_err(XrpcError::from), 453 [object] => register_repo(store, home, *object, registration, signer, now) 454 .map(|_| ()) 455 .map_err(XrpcError::from), 456 many => Err(XrpcError::internal(format!( 457 "{} repo registry objects share meta-repo", 458 many.len() 459 ))), 460 } 461} 462 463pub(crate) async fn delete_repo<H: HttpTransport, C: Clock>( 464 State(state): State<Arc<XrpcState<H, C>>>, 465 headers: HeaderMap, 466 method: crate::Method, 467 body: Bytes, 468) -> Result<Response, XrpcError> { 469 let actor = state.authenticate(&headers, &method).await?; 470 let DeleteInput { 471 did, rkey, force, .. 472 } = decode(&body)?; 473 474 let repo_did = match state.index.resolve_repo(&did, &rkey) { 475 Resolved::Ready(Some(repo_did)) => repo_did, 476 Resolved::Ready(None) => return Ok(ok_empty()), 477 Resolved::Warming => { 478 return Err(XrpcError::warming("registry projection is still warming")); 479 } 480 }; 481 482 let acl = KnotAcl::new(&state.admins, state.admission, &state.index); 483 if !can_delete_repo(&acl, &actor, &repo_did).is_allowed() { 484 return Err(XrpcError::forbidden( 485 "only repository owner or a knot admin may delete it", 486 )); 487 } 488 489 if force { 490 if !can_admin_knot(&acl, &actor).is_allowed() { 491 return Err(XrpcError::forbidden( 492 "only knot admin may force a delete past the PDS record check", 493 )); 494 } 495 } else { 496 let owner = AccountDid::from(did.clone()); 497 match state.atproto.repo_record_present(&owner, &rkey).await { 498 Ok(RecordPresence::Present) => { 499 return Err(XrpcError::conflict( 500 "sh.tangled.repo record still exists on the owner's PDS. Remove it there first or force the delete.", 501 )); 502 } 503 Ok(RecordPresence::Absent) => {} 504 Err(error) => { 505 tracing::warn!( 506 repo = %repo_did, 507 %error, 508 "proceeding w/ best-effort delete despite unconfirmed owner PDS record :p" 509 ); 510 } 511 } 512 } 513 514 let now = state.now(); 515 let knot_signer = state.secrets.signer(&state.knot_did)?; 516 let layout = state.layout.clone(); 517 let index = Arc::clone(&state.index); 518 let cob_locks = Arc::clone(&state.cob_locks); 519 let meta_path = state.meta_path.clone(); 520 let target = RepoRef { owner: did, rkey }; 521 let deleted = repo_did.clone(); 522 let home = CobHome::from(&state.knot_did); 523 let lfs_store = state.lfs.as_ref().map(|lfs| Arc::clone(&lfs.handle.store)); 524 run_blocking(move || { 525 let _repo_guard = cob_locks.repo(&deleted); 526 let _meta_guard = cob_locks.meta(); 527 let meta = Repo::open(&meta_path)?; 528 let store = CobStore::new(&meta); 529 deregister(&store, &home, target, &deleted, &knot_signer, now)?; 530 let removal = layout.remove(&deleted); 531 index.refresh_registry().map_err(|error| { 532 XrpcError::internal(format!( 533 "repo {deleted} was deregistered but registry projection refresh failed: {error}" 534 )) 535 })?; 536 if let Some(store) = &lfs_store 537 && let Err(error) = store.remove_repo(&deleted) 538 { 539 tracing::warn!( 540 repo = %deleted, 541 %error, 542 "lfs prefix removal on delete failed, orphan sweep will reclaim" 543 ); 544 } 545 removal.map_err(|error| { 546 XrpcError::internal(format!( 547 "repo {deleted} was deregistered but its on-disk directory couldn't be removed: {error}" 548 )) 549 }) 550 }) 551 .await?; 552 553 Ok(ok_empty()) 554} 555 556fn deregister( 557 store: &CobStore, 558 home: &CobHome, 559 target: RepoRef, 560 expected: &RepoDid, 561 signer: &dyn Signer, 562 now: UnixSeconds, 563) -> Result<(), XrpcError> { 564 match store 565 .list::<RepoRegistryCob>() 566 .map_err(XrpcError::from)? 567 .as_slice() 568 { 569 [] => Ok(()), 570 [object] => deregister_repo(store, home, *object, target, expected.clone(), signer, now) 571 .map(|_| ()) 572 .map_err(XrpcError::from), 573 many => Err(XrpcError::internal(format!( 574 "{} repo registry objects share meta-repo", 575 many.len() 576 ))), 577 } 578} 579 580pub(crate) async fn rename_repo<H: HttpTransport, C: Clock>( 581 State(state): State<Arc<XrpcState<H, C>>>, 582 headers: HeaderMap, 583 method: crate::Method, 584 body: Bytes, 585) -> Result<Response, XrpcError> { 586 let actor = state.authenticate(&headers, &method).await?; 587 let RenameInput { repo, rkey, name } = decode(&body)?; 588 589 let owner = match state.index.owner_of(&repo) { 590 Resolved::Ready(Some(owner)) => owner, 591 Resolved::Ready(None) => { 592 return Err(XrpcError::not_found("no such repository on this knot")); 593 } 594 Resolved::Warming => { 595 return Err(XrpcError::warming("registry projection is still warming")); 596 } 597 }; 598 599 crate::authorize_push( 600 &state, 601 &actor, 602 &repo, 603 "only repository owner or a collaborator may rename it", 604 ) 605 .await?; 606 607 let now = state.now(); 608 let knot_signer = state.secrets.signer(&state.knot_did)?; 609 let index = Arc::clone(&state.index); 610 let cob_locks = Arc::clone(&state.cob_locks); 611 let meta_path = state.meta_path.clone(); 612 let renamed = repo.clone(); 613 let home = CobHome::from(&state.knot_did); 614 run_blocking(move || { 615 { 616 let _guard = cob_locks.meta(); 617 let meta = Repo::open(&meta_path)?; 618 let store = CobStore::new(&meta); 619 match store 620 .list::<RepoRegistryCob>() 621 .map_err(XrpcError::from)? 622 .as_slice() 623 { 624 [] => { 625 return Err(XrpcError::not_found( 626 "no repositories are registered on this knot", 627 )); 628 } 629 [object] => { 630 knot_cobs::rename_repo( 631 &store, 632 &home, 633 *object, 634 Rename { 635 owner, 636 rkey, 637 name, 638 repo: renamed.clone(), 639 }, 640 &knot_signer, 641 now, 642 ) 643 .map(|_| ()) 644 .map_err(XrpcError::from)?; 645 } 646 many => { 647 return Err(XrpcError::internal(format!( 648 "{} repo registry objects share meta-repo", 649 many.len() 650 ))); 651 } 652 } 653 } 654 index.refresh_registry().map_err(|error| { 655 XrpcError::internal(format!( 656 "repo {renamed} was renamed but registry projection refresh failed: {error}" 657 )) 658 }) 659 }) 660 .await?; 661 662 Ok(ok_empty()) 663}