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 / tests.rs
103 kB 3193 lines
1use std::collections::{BTreeSet, HashMap, HashSet}; 2use std::path::PathBuf; 3use std::sync::atomic::{AtomicU64, Ordering}; 4use std::sync::{Arc, Mutex}; 5 6use axum::body::Bytes; 7use axum::extract::State; 8use axum::response::{IntoResponse, Response}; 9use base64::Engine; 10use base64::engine::general_purpose::URL_SAFE_NO_PAD; 11use futures::StreamExt; 12use http::{HeaderMap, HeaderValue, StatusCode, header::AUTHORIZATION}; 13use serde_json::json; 14use tempfile::TempDir; 15 16use knot_atproto::Atproto; 17use knot_git::Layout; 18use knot_index::{Index, Resolved}; 19use knot_runtime::{ 20 FakeHttp, HttpRequest, HttpResponse, HttpTransport, K256Signer, ManualClock, NetworkError, 21 OsEntropy, SeededEntropy, Signer, UnixMicros, 22}; 23use knot_secrets::{MasterKey, SealedStore}; 24use knot_types::{ 25 AccountDid, AdmissionPolicy, AuthorName, Email, KnotHostname, KnotId, OwnerDid, RepoDid, 26 RepoRkey, 27}; 28 29use crate::XrpcState; 30 31const KNOT_HOST: &str = "knot.nel.pet"; 32const ADMIN_HOST: &str = "admin.nel.pet"; 33const MEMBER_HOST: &str = "member.nel.pet"; 34const STRANGER_HOST: &str = "stranger.nel.pet"; 35 36const ADD_MEMBER: &str = "sh.tangled.knot.addMember"; 37const REMOVE_MEMBER: &str = "sh.tangled.knot.removeMember"; 38const BAN: &str = "sh.tangled.knot.ban"; 39const UNBAN: &str = "sh.tangled.knot.unban"; 40const CREATE: &str = "sh.tangled.repo.create"; 41const RESERVE: &str = "sh.tangled.repo.reserveKey"; 42const DELETE: &str = "sh.tangled.repo.delete"; 43const RENAME: &str = "sh.tangled.repo.rename"; 44const ADD_COLLAB: &str = "sh.tangled.repo.addCollaborator"; 45const REMOVE_COLLAB: &str = "sh.tangled.repo.removeCollaborator"; 46const SET_DEFAULT: &str = "sh.tangled.repo.setDefaultBranch"; 47const DELETE_BRANCH: &str = "sh.tangled.repo.deleteBranch"; 48 49type Responder = Box<dyn Fn(&HttpRequest) -> Result<HttpResponse, NetworkError> + Send + Sync>; 50type SharedState = Arc<XrpcState<FakeHttp<Responder>, ManualClock>>; 51 52static JTI: AtomicU64 = AtomicU64::new(0); 53 54fn knot_did() -> KnotId { 55 KnotId::new(format!("did:web:{KNOT_HOST}")).unwrap() 56} 57 58fn account(host: &str) -> AccountDid { 59 AccountDid::new(format!("did:web:{host}")).unwrap() 60} 61 62fn signer(seed: u64) -> K256Signer { 63 K256Signer::generate(&SeededEntropy::new(seed)) 64} 65 66fn did_web_doc(did: &str, sec1: &[u8], pds: &str) -> Bytes { 67 let multikey = knot_types::crypto::multikey(0xe7, sec1); 68 Bytes::from( 69 serde_json::to_vec(&json!({ 70 "id": did, 71 "alsoKnownAs": [], 72 "verificationMethod": [{ 73 "id": format!("{did}#atproto"), 74 "type": "Multikey", 75 "controller": did, 76 "publicKeyMultibase": multikey 77 }], 78 "service": [{ 79 "id": "#atproto_pds", 80 "type": "AtprotoPersonalDataServer", 81 "serviceEndpoint": pds 82 }] 83 })) 84 .unwrap(), 85 ) 86} 87 88fn repo_did_doc(did: &str, multikey: &str) -> Bytes { 89 Bytes::from( 90 serde_json::to_vec(&json!({ 91 "id": did, 92 "verificationMethod": [{ 93 "id": format!("{did}#repo"), 94 "type": "Multikey", 95 "controller": did, 96 "publicKeyMultibase": multikey 97 }] 98 })) 99 .unwrap(), 100 ) 101} 102 103fn mint(signer: &K256Signer, issuer: &AccountDid, method: &str) -> String { 104 let jti = JTI.fetch_add(1, Ordering::Relaxed); 105 let header = URL_SAFE_NO_PAD.encode(br#"{"alg":"ES256K","typ":"JWT"}"#); 106 let payload = URL_SAFE_NO_PAD.encode( 107 serde_json::to_vec(&json!({ 108 "iss": issuer.as_str(), 109 "aud": format!("did:web:{KNOT_HOST}"), 110 "exp": 1_100, 111 "iat": 999, 112 "jti": format!("nonce-{jti}"), 113 "lxm": method, 114 })) 115 .unwrap(), 116 ); 117 let signing_input = format!("{header}.{payload}"); 118 let signature = signer.sign(signing_input.as_bytes()); 119 format!( 120 "{signing_input}.{}", 121 URL_SAFE_NO_PAD.encode(signature.as_bytes()) 122 ) 123} 124 125fn bearer(token: &str) -> HeaderMap { 126 let mut headers = HeaderMap::new(); 127 headers.insert( 128 AUTHORIZATION, 129 HeaderValue::from_str(&format!("Bearer {token}")).unwrap(), 130 ); 131 headers 132} 133 134fn body(value: serde_json::Value) -> Bytes { 135 Bytes::from(serde_json::to_vec(&value).unwrap()) 136} 137 138fn into_response(result: Result<Response, crate::XrpcError>) -> Response { 139 match result { 140 Ok(response) => response, 141 Err(error) => error.into_response(), 142 } 143} 144 145async fn json_of(response: Response) -> serde_json::Value { 146 let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) 147 .await 148 .unwrap(); 149 serde_json::from_slice(&bytes).unwrap() 150} 151 152async fn call<F, Fut>( 153 world: &World, 154 handler: F, 155 signer: &K256Signer, 156 host: &str, 157 nsid: &str, 158 value: serde_json::Value, 159) -> Response 160where 161 F: FnOnce(State<SharedState>, HeaderMap, crate::Method, Bytes) -> Fut, 162 Fut: std::future::Future<Output = Result<Response, crate::XrpcError>>, 163{ 164 let token = mint(signer, &account(host), nsid); 165 into_response( 166 handler( 167 world.state(), 168 bearer(&token), 169 crate::Method::from_nsid(nsid), 170 body(value), 171 ) 172 .await, 173 ) 174} 175 176async fn as_member<F, Fut>( 177 world: &World, 178 handler: F, 179 nsid: &str, 180 value: serde_json::Value, 181) -> Response 182where 183 F: FnOnce(State<SharedState>, HeaderMap, crate::Method, Bytes) -> Fut, 184 Fut: std::future::Future<Output = Result<Response, crate::XrpcError>>, 185{ 186 call(world, handler, &world.member, MEMBER_HOST, nsid, value).await 187} 188 189async fn as_admin<F, Fut>( 190 world: &World, 191 handler: F, 192 nsid: &str, 193 value: serde_json::Value, 194) -> Response 195where 196 F: FnOnce(State<SharedState>, HeaderMap, crate::Method, Bytes) -> Fut, 197 Fut: std::future::Future<Output = Result<Response, crate::XrpcError>>, 198{ 199 call(world, handler, &world.admin, ADMIN_HOST, nsid, value).await 200} 201 202async fn as_stranger<F, Fut>( 203 world: &World, 204 handler: F, 205 nsid: &str, 206 value: serde_json::Value, 207) -> Response 208where 209 F: FnOnce(State<SharedState>, HeaderMap, crate::Method, Bytes) -> Fut, 210 Fut: std::future::Future<Output = Result<Response, crate::XrpcError>>, 211{ 212 call(world, handler, &world.stranger, STRANGER_HOST, nsid, value).await 213} 214 215fn member_owner() -> OwnerDid { 216 OwnerDid::new(format!("did:web:{MEMBER_HOST}")).unwrap() 217} 218 219fn resolve(world: &World, rkey: &str) -> Resolved<Option<RepoDid>> { 220 world 221 .state 222 .index 223 .resolve_repo(&member_owner(), &RepoRkey::new(rkey).unwrap()) 224} 225 226fn replay(world: &World) -> Vec<std::sync::Arc<knot_events::Event>> { 227 world 228 .state 229 .events 230 .replay( 231 knot_events::EventCursor::START, 232 knot_events::ReplayBounds::new( 233 knot_events::ReplayEvents::new(64).unwrap(), 234 knot_events::ReplayBytes::new(16 << 20).unwrap(), 235 ), 236 ) 237 .events 238} 239 240fn event_count(world: &World) -> usize { 241 replay(world).len() 242} 243 244fn last_event(world: &World, nsid: &str) -> std::sync::Arc<knot_events::Event> { 245 replay(world) 246 .into_iter() 247 .rev() 248 .find(|event| event.nsid == nsid) 249 .unwrap_or_else(|| panic!("{nsid} event is emitted")) 250} 251 252fn git_events(world: &World) -> Vec<std::sync::Arc<knot_events::Event>> { 253 replay(world) 254 .into_iter() 255 .filter(|event| { 256 !matches!( 257 event.nsid, 258 "sh.tangled.knot.memberUpdate" | "sh.tangled.repo.collaboratorUpdate" 259 ) 260 }) 261 .collect() 262} 263 264fn only_git_event(world: &World) -> std::sync::Arc<knot_events::Event> { 265 let mut events = git_events(world); 266 assert_eq!(events.len(), 1, "expected exactly one non-acl event"); 267 events.remove(0) 268} 269 270fn bootstrap(dir: &TempDir, rebuild: bool) -> (Layout, Arc<Index>, PathBuf) { 271 let scan_path = dir.path().join("repos"); 272 std::fs::create_dir_all(&scan_path).unwrap(); 273 let knot = knot_did(); 274 let layout = Layout::new(&scan_path).reserving_meta(&knot).unwrap(); 275 layout.bootstrap_meta(&knot).unwrap(); 276 let meta_path = layout.meta_path(&knot).unwrap(); 277 let index = Arc::new(Index::new(meta_path.clone(), layout.clone())); 278 if rebuild { 279 index.rebuild().unwrap(); 280 } 281 (layout, index, meta_path) 282} 283 284fn state_from( 285 dir: &TempDir, 286 boot: (Layout, Arc<Index>, PathBuf), 287 responder: Responder, 288 admission: AdmissionPolicy, 289 reservations: Arc<crate::Reservations>, 290 git_http: Arc<dyn HttpTransport>, 291) -> SharedState { 292 let (layout, index, meta_path) = boot; 293 let knot = knot_did(); 294 let knot_url = knot_types::KnotServiceUrl::new(format!("https://{KNOT_HOST}")).unwrap(); 295 let atproto = Arc::new(Atproto::new( 296 FakeHttp::new(responder), 297 ManualClock::new(UnixMicros::new(1_000_000_000)), 298 knot.clone(), 299 knot_atproto::PlcDirectory::new(url::Url::parse("https://plc.directory/").unwrap()) 300 .unwrap(), 301 )); 302 let secrets = Arc::new( 303 SealedStore::open( 304 dir.path().join("keys.sealed"), 305 &MasterKey::new([7u8; 32]).unwrap(), 306 Box::new(OsEntropy), 307 ) 308 .unwrap(), 309 ); 310 secrets.ensure(&knot).unwrap(); 311 Arc::new(XrpcState { 312 layout, 313 index, 314 atproto, 315 secrets, 316 entropy: Arc::new(OsEntropy), 317 ci_logs: None, 318 admins: BTreeSet::from([account(ADMIN_HOST)]), 319 admission, 320 knot_did: knot, 321 knot_hostname: KnotHostname::new(KNOT_HOST).unwrap(), 322 meta_path, 323 knot_service_url: knot_url, 324 limiter: Arc::new(crate::PreAuthLimiter::default()), 325 cob_locks: Arc::new(crate::CobLocks::default()), 326 reservations, 327 trusted_proxy_header: None, 328 committer: crate::Committer { 329 name: AuthorName::new("Tangled"), 330 email: Email::new("noreply@tangled.sh"), 331 }, 332 byte_limits: crate::ByteLimits::default(), 333 budgets: crate::Budgets::default(), 334 git_http, 335 pack_limits: knot_pack::PackLimits::default(), 336 service_owner: account(ADMIN_HOST), 337 events: Arc::new(knot_events::EventLog::new( 338 ManualClock::new(UnixMicros::new(1_000_000_000)), 339 knot_events::ReplayBounds::new( 340 knot_events::ReplayEvents::new(1024).unwrap(), 341 knot_events::ReplayBytes::new(16 << 20).unwrap(), 342 ), 343 )), 344 subscriber_gate: Arc::new(knot_events::SubscriberGate::new( 345 knot_events::GlobalSubscriberLimit::new(16), 346 knot_events::PerPeerSubscriberLimit::new(8), 347 )), 348 maintenance: knot_maintenance::MaintenanceHandle::disabled(), 349 appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), 350 slots: knot_resource::Slots::testing(8), 351 lfs: None, 352 catalog: Arc::new(knot_messages::Catalog::defaults()), 353 }) 354} 355 356fn world_responder( 357 pubkeys: HashMap<String, Vec<u8>>, 358 repo_docs: Arc<Mutex<HashMap<String, String>>>, 359 pds_records: Arc<Mutex<HashSet<String>>>, 360) -> Responder { 361 let doc_url = format!("https://{KNOT_HOST}"); 362 Box::new(move |request: &HttpRequest| { 363 if request.method == http::Method::POST { 364 return Ok(HttpResponse { 365 status: StatusCode::OK, 366 headers: http::HeaderMap::new(), 367 body: Bytes::new(), 368 }); 369 } 370 if request.url.path().ends_with("com.atproto.repo.getRecord") { 371 let rkey = request 372 .url 373 .query_pairs() 374 .find(|(key, _)| key == "rkey") 375 .map(|(_, value)| value.into_owned()) 376 .unwrap_or_default(); 377 let present = pds_records.lock().unwrap().contains(&rkey); 378 return Ok(HttpResponse { 379 status: if present { 380 StatusCode::OK 381 } else { 382 StatusCode::BAD_REQUEST 383 }, 384 headers: http::HeaderMap::new(), 385 body: if present { 386 Bytes::new() 387 } else { 388 Bytes::from_static(b"{\"error\":\"RecordNotFound\"}") 389 }, 390 }); 391 } 392 let host = request.url.host_str().unwrap_or_default(); 393 if let Some(multikey) = repo_docs.lock().unwrap().get(host).cloned() { 394 return Ok(HttpResponse { 395 status: StatusCode::OK, 396 headers: http::HeaderMap::new(), 397 body: repo_did_doc(&format!("did:web:{host}"), &multikey), 398 }); 399 } 400 match pubkeys.get(host) { 401 Some(sec1) => Ok(HttpResponse { 402 status: StatusCode::OK, 403 headers: http::HeaderMap::new(), 404 body: did_web_doc(&format!("did:web:{host}"), sec1, &doc_url), 405 }), 406 None => Ok(HttpResponse { 407 status: StatusCode::NOT_FOUND, 408 headers: http::HeaderMap::new(), 409 body: Bytes::new(), 410 }), 411 } 412 }) 413} 414 415struct World { 416 _dir: TempDir, 417 layout: Layout, 418 state: SharedState, 419 admin: K256Signer, 420 member: K256Signer, 421 stranger: K256Signer, 422 repo_docs: Arc<Mutex<HashMap<String, String>>>, 423 pds_records: Arc<Mutex<HashSet<String>>>, 424} 425 426impl World { 427 fn new() -> Self { 428 Self::build(256, 256, AdmissionPolicy::Closed, None) 429 } 430 431 fn open() -> Self { 432 Self::build(256, 256, AdmissionPolicy::Open, None) 433 } 434 435 fn with_pending_limit(limit: usize) -> Self { 436 Self::build(limit, limit, AdmissionPolicy::Closed, None) 437 } 438 439 fn with_limits(global: usize, per_actor: usize) -> Self { 440 Self::build(global, per_actor, AdmissionPolicy::Closed, None) 441 } 442 443 fn with_git_http(git_http: Arc<dyn HttpTransport>) -> Self { 444 Self::build(256, 256, AdmissionPolicy::Closed, Some(git_http)) 445 } 446 447 fn build( 448 global: usize, 449 per_actor: usize, 450 admission: AdmissionPolicy, 451 git_http: Option<Arc<dyn HttpTransport>>, 452 ) -> Self { 453 let dir = tempfile::tempdir().unwrap(); 454 let (layout, index, meta_path) = bootstrap(&dir, true); 455 let admin = signer(1); 456 let member = signer(2); 457 let stranger = signer(3); 458 let repo_docs: Arc<Mutex<HashMap<String, String>>> = Arc::new(Mutex::new(HashMap::new())); 459 let pds_records: Arc<Mutex<HashSet<String>>> = Arc::new(Mutex::new(HashSet::new())); 460 let pubkeys = HashMap::from([ 461 ( 462 ADMIN_HOST.to_string(), 463 admin.public_key().as_bytes().to_vec(), 464 ), 465 ( 466 MEMBER_HOST.to_string(), 467 member.public_key().as_bytes().to_vec(), 468 ), 469 ( 470 STRANGER_HOST.to_string(), 471 stranger.public_key().as_bytes().to_vec(), 472 ), 473 ]); 474 let responder = world_responder(pubkeys, Arc::clone(&repo_docs), Arc::clone(&pds_records)); 475 let state = state_from( 476 &dir, 477 (layout.clone(), index, meta_path), 478 responder, 479 admission, 480 Arc::new(crate::Reservations::new( 481 crate::ReservationTtl::new(1_000_000), 482 crate::PerActorQuota::new(per_actor), 483 crate::GlobalQuota::new(global), 484 )), 485 git_http.unwrap_or_else(no_git_upstream), 486 ); 487 Self { 488 _dir: dir, 489 layout, 490 state, 491 admin, 492 member, 493 stranger, 494 repo_docs, 495 pds_records, 496 } 497 } 498 499 fn state(&self) -> State<SharedState> { 500 State(Arc::clone(&self.state)) 501 } 502 503 fn publish_repo_doc(&self, host: &str, multikey: &str) { 504 self.repo_docs 505 .lock() 506 .unwrap() 507 .insert(host.to_string(), multikey.to_string()); 508 } 509 510 fn publish_pds_record(&self, rkey: &str) { 511 self.pds_records.lock().unwrap().insert(rkey.to_string()); 512 } 513} 514 515fn build_state(responder: Responder, rebuild: bool) -> (TempDir, SharedState) { 516 let dir = tempfile::tempdir().unwrap(); 517 let (layout, index, meta_path) = bootstrap(&dir, rebuild); 518 let state = state_from( 519 &dir, 520 (layout, index, meta_path), 521 responder, 522 AdmissionPolicy::Closed, 523 Arc::new(crate::Reservations::new( 524 crate::ReservationTtl::new(1_000_000), 525 crate::PerActorQuota::new(256), 526 crate::GlobalQuota::new(256), 527 )), 528 no_git_upstream(), 529 ); 530 (dir, state) 531} 532 533fn no_git_upstream() -> Arc<dyn HttpTransport> { 534 Arc::new(FakeHttp::new(|_request: &HttpRequest| { 535 Err(NetworkError::Connect( 536 "no git upstream is served in this test".to_string(), 537 )) 538 })) 539} 540 541fn doc_responder( 542 sec1: Vec<u8>, 543 post_status: impl Fn() -> StatusCode + Send + Sync + 'static, 544) -> Responder { 545 let pds = format!("https://{KNOT_HOST}"); 546 Box::new(move |request: &HttpRequest| { 547 let post = request.method == http::Method::POST; 548 let host = request.url.host_str().unwrap_or_default(); 549 Ok(HttpResponse { 550 status: if post { post_status() } else { StatusCode::OK }, 551 headers: http::HeaderMap::new(), 552 body: if post { 553 Bytes::new() 554 } else { 555 did_web_doc(&format!("did:web:{host}"), &sec1, &pds) 556 }, 557 }) 558 }) 559} 560 561async fn add_member_helper(world: &World) { 562 assert_eq!( 563 as_admin( 564 world, 565 crate::members::add_member, 566 ADD_MEMBER, 567 json!({ "subject": format!("did:web:{MEMBER_HOST}") }) 568 ) 569 .await 570 .status(), 571 StatusCode::OK 572 ); 573} 574 575async fn create_repo_helper(world: &World, name: &str) -> RepoDid { 576 assert_eq!( 577 as_member( 578 world, 579 crate::repos::create_repo, 580 CREATE, 581 json!({ "rkey": name, "name": name }) 582 ) 583 .await 584 .status(), 585 StatusCode::OK 586 ); 587 match resolve(world, name) { 588 Resolved::Ready(Some(did)) => did, 589 other => panic!("repo {name} wasn't registered: {other:?}"), 590 } 591} 592 593async fn create_status( 594 world: &World, 595 signer: &K256Signer, 596 host: &str, 597 value: serde_json::Value, 598) -> StatusCode { 599 call( 600 world, 601 crate::repos::create_repo, 602 signer, 603 host, 604 CREATE, 605 value, 606 ) 607 .await 608 .status() 609} 610 611async fn reserve_status(world: &World, signer: &K256Signer, host: &str, did: &str) -> StatusCode { 612 call( 613 world, 614 crate::repos::reserve_key, 615 signer, 616 host, 617 RESERVE, 618 json!({ "repoDid": did }), 619 ) 620 .await 621 .status() 622} 623 624async fn reserve_repo_key(world: &World, did_web: &str) -> String { 625 let response = as_member( 626 world, 627 crate::repos::reserve_key, 628 RESERVE, 629 json!({ "repoDid": did_web }), 630 ) 631 .await; 632 assert_eq!(response.status(), StatusCode::OK); 633 let key = json_of(response).await["key"].as_str().unwrap().to_string(); 634 let host = did_web.strip_prefix("did:web:").unwrap(); 635 world.publish_repo_doc(host, &key); 636 key 637} 638 639async fn rename_repo_as( 640 world: &World, 641 signer: &K256Signer, 642 host: &str, 643 repo: &RepoDid, 644 rkey: &str, 645) -> StatusCode { 646 call( 647 world, 648 crate::repos::rename_repo, 649 signer, 650 host, 651 RENAME, 652 json!({ "repo": repo.as_str(), "rkey": rkey, "name": rkey }), 653 ) 654 .await 655 .status() 656} 657 658#[tokio::test] 659async fn admission_gates_repo_creation() { 660 let closed = World::new(); 661 let make = || json!({ "rkey": "anemone", "name": "anemone" }); 662 assert_eq!( 663 create_status(&closed, &closed.stranger, STRANGER_HOST, make()).await, 664 StatusCode::FORBIDDEN, 665 "a closed knot denies a stranger" 666 ); 667 668 let open = World::open(); 669 assert_eq!( 670 create_status(&open, &open.stranger, STRANGER_HOST, make()).await, 671 StatusCode::OK 672 ); 673 assert!( 674 matches!( 675 open.state.index.resolve_repo( 676 &OwnerDid::new(format!("did:web:{STRANGER_HOST}")).unwrap(), 677 &RepoRkey::new("anemone").unwrap() 678 ), 679 Resolved::Ready(Some(_)) 680 ), 681 "an open knot registers the stranger's repo without membership" 682 ); 683} 684 685#[tokio::test] 686async fn blocklist_lifecycle() { 687 let world = World::open(); 688 let subject = format!("did:web:{STRANGER_HOST}"); 689 let make = || json!({ "rkey": "anemone", "name": "anemone" }); 690 691 assert_eq!( 692 as_admin( 693 &world, 694 crate::blocklist::ban, 695 BAN, 696 json!({ "subject": subject }) 697 ) 698 .await 699 .status(), 700 StatusCode::OK 701 ); 702 assert!(matches!( 703 world.state.index.is_blocked(&account(STRANGER_HOST)), 704 Resolved::Ready(true) 705 )); 706 assert_eq!( 707 create_status(&world, &world.stranger, STRANGER_HOST, make()).await, 708 StatusCode::FORBIDDEN, 709 "a banned account cannot create" 710 ); 711 712 assert_eq!( 713 as_admin( 714 &world, 715 crate::blocklist::unban, 716 UNBAN, 717 json!({ "subject": subject }) 718 ) 719 .await 720 .status(), 721 StatusCode::OK 722 ); 723 assert!(matches!( 724 world.state.index.is_blocked(&account(STRANGER_HOST)), 725 Resolved::Ready(false) 726 )); 727 assert_eq!( 728 create_status(&world, &world.stranger, STRANGER_HOST, make()).await, 729 StatusCode::OK, 730 "an unban restores creation" 731 ); 732 733 assert_eq!( 734 as_admin( 735 &world, 736 crate::blocklist::ban, 737 BAN, 738 json!({ "subject": format!("did:web:{ADMIN_HOST}") }) 739 ) 740 .await 741 .status(), 742 StatusCode::FORBIDDEN, 743 "an admin cannot be banned" 744 ); 745 assert_eq!( 746 as_stranger( 747 &world, 748 crate::blocklist::ban, 749 BAN, 750 json!({ "subject": format!("did:web:{MEMBER_HOST}") }) 751 ) 752 .await 753 .status(), 754 StatusCode::FORBIDDEN, 755 "a non-admin cannot ban" 756 ); 757} 758 759#[tokio::test] 760async fn member_lifecycle() { 761 let world = World::new(); 762 let subject = json!({ "subject": format!("did:web:{MEMBER_HOST}") }); 763 764 assert_eq!( 765 as_admin( 766 &world, 767 crate::members::add_member, 768 ADD_MEMBER, 769 subject.clone() 770 ) 771 .await 772 .status(), 773 StatusCode::OK 774 ); 775 assert_eq!( 776 world.state.index.is_member(&account(MEMBER_HOST)), 777 Resolved::Ready(true), 778 "member is effective on the very next read, with no firehose" 779 ); 780 let events = replay(&world); 781 let [added] = events.as_slice() else { 782 panic!("expected exactly one event, got {}", events.len()); 783 }; 784 assert_eq!(added.nsid, "sh.tangled.knot.memberUpdate"); 785 assert_eq!(added.payload["op"], "add"); 786 assert_eq!(added.payload["subject"], account(MEMBER_HOST).to_string()); 787 788 let baseline = event_count(&world); 789 assert_eq!( 790 as_admin( 791 &world, 792 crate::members::add_member, 793 ADD_MEMBER, 794 subject.clone() 795 ) 796 .await 797 .status(), 798 StatusCode::OK 799 ); 800 assert_eq!( 801 event_count(&world), 802 baseline, 803 "a redundant add is a no-op and emits no event" 804 ); 805 806 assert_eq!( 807 as_admin( 808 &world, 809 crate::members::remove_member, 810 REMOVE_MEMBER, 811 subject.clone() 812 ) 813 .await 814 .status(), 815 StatusCode::OK 816 ); 817 assert_eq!( 818 world.state.index.is_member(&account(MEMBER_HOST)), 819 Resolved::Ready(false), 820 "removed member is gone on the very next read" 821 ); 822 let removed = last_event(&world, "sh.tangled.knot.memberUpdate"); 823 assert_eq!(removed.payload["op"], "remove"); 824 assert_eq!(removed.payload["subject"], account(MEMBER_HOST).to_string()); 825 826 let baseline = event_count(&world); 827 assert_eq!( 828 as_admin( 829 &world, 830 crate::members::remove_member, 831 REMOVE_MEMBER, 832 subject 833 ) 834 .await 835 .status(), 836 StatusCode::OK 837 ); 838 assert_eq!( 839 event_count(&world), 840 baseline, 841 "removing a non-member is a no-op and emits no event" 842 ); 843} 844 845#[tokio::test] 846async fn add_member_auth_outcomes() { 847 struct Case { 848 headers: HeaderMap, 849 body: serde_json::Value, 850 status: StatusCode, 851 why: &'static str, 852 } 853 854 let world = World::new(); 855 let lowercase = { 856 let token = mint(&world.admin, &account(ADMIN_HOST), ADD_MEMBER); 857 let mut headers = HeaderMap::new(); 858 headers.insert( 859 AUTHORIZATION, 860 HeaderValue::from_str(&format!("bearer {token}")).unwrap(), 861 ); 862 headers 863 }; 864 let cases = vec![ 865 Case { 866 headers: bearer(&mint(&world.member, &account(MEMBER_HOST), ADD_MEMBER)), 867 body: json!({ "subject": "did:web:olaren.dev" }), 868 status: StatusCode::FORBIDDEN, 869 why: "a non-admin cannot add a member", 870 }, 871 Case { 872 headers: HeaderMap::new(), 873 body: json!({ "subject": format!("did:web:{MEMBER_HOST}") }), 874 status: StatusCode::UNAUTHORIZED, 875 why: "a request without a token is unauthorized", 876 }, 877 Case { 878 headers: bearer(&mint(&world.admin, &account(ADMIN_HOST), ADD_MEMBER)), 879 body: json!({ "subject": "not-a-did" }), 880 status: StatusCode::BAD_REQUEST, 881 why: "an invalid DID is rejected at decode by the newtype Deserialize, never reaching a handler", 882 }, 883 Case { 884 headers: lowercase, 885 body: json!({ "subject": "did:web:olaren.dev" }), 886 status: StatusCode::OK, 887 why: "the bearer scheme is matched case-insensitively per RFC 7235", 888 }, 889 ]; 890 891 let world_ref = &world; 892 futures::stream::iter(cases) 893 .for_each(|case| async move { 894 let status = into_response( 895 crate::members::add_member( 896 world_ref.state(), 897 case.headers, 898 crate::Method::from_nsid(ADD_MEMBER), 899 body(case.body), 900 ) 901 .await, 902 ) 903 .status(); 904 assert_eq!(status, case.status, "{}", case.why); 905 }) 906 .await; 907} 908 909#[tokio::test] 910async fn a_transient_upstream_identity_failure_is_a_503_named_upstream_unavailable() { 911 let responder: Responder = Box::new(|_request: &HttpRequest| { 912 Ok(HttpResponse { 913 status: StatusCode::INTERNAL_SERVER_ERROR, 914 headers: http::HeaderMap::new(), 915 body: Bytes::new(), 916 }) 917 }); 918 let (_dir, state) = build_state(responder, true); 919 let admin = signer(1); 920 let token = mint(&admin, &account(ADMIN_HOST), ADD_MEMBER); 921 let response = into_response( 922 crate::members::add_member( 923 State(state), 924 bearer(&token), 925 crate::Method::from_nsid(ADD_MEMBER), 926 body(json!({ "subject": format!("did:web:{MEMBER_HOST}") })), 927 ) 928 .await, 929 ); 930 assert_eq!( 931 response.status(), 932 StatusCode::SERVICE_UNAVAILABLE, 933 "transient issuer-doc failure is 503, distinct from a bad token's 401" 934 ); 935 assert_eq!( 936 json_of(response).await["error"], 937 "UpstreamUnavailable", 938 "transient upstream identity failure is named distinctly from a warming projection" 939 ); 940} 941 942#[tokio::test] 943async fn concurrent_first_member_adds_converge_to_a_single_cob_object() { 944 use knot_cob::CobStore; 945 use knot_cobs::MembersCob; 946 use knot_git::Repo; 947 948 let world = World::new(); 949 let (left, right) = tokio::join!( 950 as_admin( 951 &world, 952 crate::members::add_member, 953 ADD_MEMBER, 954 json!({ "subject": "did:web:witchcraft.systems" }) 955 ), 956 as_admin( 957 &world, 958 crate::members::add_member, 959 ADD_MEMBER, 960 json!({ "subject": "did:web:isabelroses.com" }) 961 ), 962 ); 963 assert_eq!(left.status(), StatusCode::OK); 964 assert_eq!(right.status(), StatusCode::OK); 965 966 let meta = Repo::open(world.state.meta_path.clone()).unwrap(); 967 assert_eq!( 968 CobStore::new(&meta).list::<MembersCob>().unwrap().len(), 969 1, 970 "concurrent first adds serialize onto one singleton members COB, never splitting it" 971 ); 972 assert_eq!( 973 world.state.index.is_member(&account("witchcraft.systems")), 974 Resolved::Ready(true) 975 ); 976 assert_eq!( 977 world.state.index.is_member(&account("isabelroses.com")), 978 Resolved::Ready(true) 979 ); 980} 981 982#[tokio::test] 983async fn re_adding_a_member_under_warming_appends_no_redundant_change() { 984 use knot_cob::CobStore; 985 use knot_cobs::{Grant, MembersChange, MembersCob}; 986 use knot_git::Repo; 987 988 let admin = signer(1); 989 let responder = doc_responder(admin.public_key().as_bytes().to_vec(), || StatusCode::OK); 990 let (_dir, state) = build_state(responder, false); 991 992 let now = state.now(); 993 let knot_signer = state.secrets.signer(&state.knot_did).unwrap(); 994 let subject = account(MEMBER_HOST); 995 let meta = Repo::open(state.meta_path.clone()).unwrap(); 996 let created = CobStore::new(&meta) 997 .create( 998 &knot_cob::CobHome::from(&state.knot_did), 999 &MembersChange::Add(Grant { 1000 subject: subject.clone(), 1001 added_by: account(ADMIN_HOST), 1002 created_at: now, 1003 }), 1004 &knot_signer, 1005 now, 1006 ) 1007 .unwrap(); 1008 1009 let token = mint(&admin, &account(ADMIN_HOST), ADD_MEMBER); 1010 assert_eq!( 1011 into_response( 1012 crate::members::add_member( 1013 State(Arc::clone(&state)), 1014 bearer(&token), 1015 crate::Method::from_nsid(ADD_MEMBER), 1016 body(json!({ "subject": subject.as_str() })), 1017 ) 1018 .await 1019 ) 1020 .status(), 1021 StatusCode::OK 1022 ); 1023 1024 let meta = Repo::open(state.meta_path.clone()).unwrap(); 1025 let delta = CobStore::new(&meta) 1026 .changes_since::<MembersCob>(created.object, Some(created.tip)) 1027 .unwrap(); 1028 assert!( 1029 delta.changes.is_empty(), 1030 "re-adding an existing member appends no change, even while projection is warming" 1031 ); 1032} 1033 1034#[tokio::test] 1035async fn create_mints_a_did_plc_repo_and_refuses_a_duplicate_name() { 1036 let world = World::new(); 1037 add_member_helper(&world).await; 1038 1039 let repo_did = create_repo_helper(&world, "anemone").await; 1040 assert!( 1041 repo_did.as_str().starts_with("did:plc:"), 1042 "knot minted a did:plc identity for the repo" 1043 ); 1044 assert!( 1045 world.layout.open(&repo_did).is_ok(), 1046 "bare repo exists on disk under its minted DID" 1047 ); 1048 assert_eq!( 1049 world.state.secrets.len(), 1050 1, 1051 "only the shared knot key is sealed" 1052 ); 1053 1054 assert_eq!( 1055 create_status( 1056 &world, 1057 &world.member, 1058 MEMBER_HOST, 1059 json!({ "rkey": "anemone", "name": "anemone" }) 1060 ) 1061 .await, 1062 StatusCode::CONFLICT, 1063 "a second repo of the same name is refused, never silently overwritten" 1064 ); 1065 assert_eq!( 1066 resolve(&world, "anemone"), 1067 Resolved::Ready(Some(repo_did.clone())), 1068 "the original repo still owns the name" 1069 ); 1070 assert!( 1071 world.layout.open(&repo_did).is_ok(), 1072 "the original repo is untouched on disk" 1073 ); 1074} 1075 1076#[tokio::test] 1077async fn create_and_reserve_reject_bad_identities() { 1078 let world = World::new(); 1079 add_member_helper(&world).await; 1080 1081 assert_eq!( 1082 create_status( 1083 &world, 1084 &world.member, 1085 MEMBER_HOST, 1086 json!({ "rkey": "a", "name": "a", "repoDid": "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa" }) 1087 ) 1088 .await, 1089 StatusCode::BAD_REQUEST, 1090 "a did:plc cannot be brought; the knot mints those itself" 1091 ); 1092 1093 assert_eq!( 1094 create_status( 1095 &world, 1096 &world.member, 1097 MEMBER_HOST, 1098 json!({ "rkey": "evil", "name": "evil", "repoDid": format!("did:web:{KNOT_HOST}") }) 1099 ) 1100 .await, 1101 StatusCode::BAD_REQUEST, 1102 "the knot's own DID as a repoDid is a client error instead of a 500" 1103 ); 1104 assert_eq!( 1105 world.state.index.is_member(&account(MEMBER_HOST)), 1106 Resolved::Ready(true), 1107 "the meta-repo is intact" 1108 ); 1109 1110 assert_eq!( 1111 create_status( 1112 &world, 1113 &world.member, 1114 MEMBER_HOST, 1115 json!({ "rkey": "uni", "name": "uni", "repoDid": "did:web:uni.olaren.dev" }) 1116 ) 1117 .await, 1118 StatusCode::BAD_REQUEST, 1119 "a did:web create without a prior reserveKey is refused" 1120 ); 1121 assert_eq!( 1122 resolve(&world, "uni"), 1123 Resolved::Ready(None), 1124 "nothing was registered for the unproven did:web" 1125 ); 1126 1127 assert_eq!( 1128 reserve_status( 1129 &world, 1130 &world.member, 1131 MEMBER_HOST, 1132 &format!("did:web:{KNOT_HOST}") 1133 ) 1134 .await, 1135 StatusCode::BAD_REQUEST, 1136 "the knot's own DID cannot have a repo key reserved against it" 1137 ); 1138 1139 let unpublished = "did:web:conch.olaren.dev"; 1140 assert_eq!( 1141 reserve_status(&world, &world.member, MEMBER_HOST, unpublished).await, 1142 StatusCode::OK 1143 ); 1144 let impostor = knot_types::crypto::multikey(0xe7, signer(99).public_key().as_bytes()); 1145 world.publish_repo_doc("conch.olaren.dev", &impostor); 1146 assert_eq!( 1147 create_status( 1148 &world, 1149 &world.member, 1150 MEMBER_HOST, 1151 json!({ "rkey": "conch", "name": "conch", "repoDid": unpublished }) 1152 ) 1153 .await, 1154 StatusCode::BAD_REQUEST, 1155 "a did:web whose document publishes a different key fails the control proof" 1156 ); 1157 assert!( 1158 world 1159 .layout 1160 .open(&RepoDid::new(unpublished).unwrap()) 1161 .is_err(), 1162 "the repo was never created on disk" 1163 ); 1164 1165 let victim = "did:web:victim.olaren.dev"; 1166 let reserved = reserve_repo_key(&world, victim).await; 1167 assert_eq!( 1168 reserve_repo_key(&world, victim).await, 1169 reserved, 1170 "re-reserving returns the same key so an already-published document stays valid" 1171 ); 1172 assert_eq!( 1173 create_status( 1174 &world, 1175 &world.admin, 1176 ADMIN_HOST, 1177 json!({ "rkey": "victim", "name": "victim", "repoDid": victim }) 1178 ) 1179 .await, 1180 StatusCode::BAD_REQUEST, 1181 "only the account that reserved the did:web may create it" 1182 ); 1183 assert!( 1184 world.layout.open(&RepoDid::new(victim).unwrap()).is_err(), 1185 "the hijack attempt created no repo on disk" 1186 ); 1187 1188 let before = world.state.secrets.len(); 1189 assert_eq!( 1190 create_status( 1191 &world, 1192 &world.member, 1193 MEMBER_HOST, 1194 json!({ "rkey": "doomed", "name": "doomed", "defaultBranch": "bad..name" }) 1195 ) 1196 .await, 1197 StatusCode::BAD_REQUEST 1198 ); 1199 assert_eq!( 1200 world.state.secrets.len(), 1201 before, 1202 "an invalid branch is rejected before any key is minted, sealed, or DID submitted" 1203 ); 1204} 1205 1206#[tokio::test] 1207async fn a_byo_did_web_repo_is_accepted_and_its_key_is_returned() { 1208 let world = World::new(); 1209 add_member_helper(&world).await; 1210 let did = "did:web:nautilus.olaren.dev"; 1211 let reserved_key = reserve_repo_key(&world, did).await; 1212 1213 let response = as_member( 1214 &world, 1215 crate::repos::create_repo, 1216 CREATE, 1217 json!({ "rkey": "nautilus", "name": "nautilus", "repoDid": did }), 1218 ) 1219 .await; 1220 assert_eq!(response.status(), StatusCode::OK); 1221 let output = json_of(response).await; 1222 assert_eq!(output["repoDid"].as_str(), Some(did)); 1223 1224 let repo_did = RepoDid::new(did).unwrap(); 1225 assert!( 1226 world.layout.open(&repo_did).is_ok(), 1227 "bring-your-own did:web repo is on disk" 1228 ); 1229 let knot_public = world 1230 .state 1231 .secrets 1232 .public_key(&world.state.knot_did) 1233 .unwrap(); 1234 let expected_key = knot_types::crypto::multikey(0xe7, knot_public.as_bytes()); 1235 assert_eq!( 1236 expected_key, reserved_key, 1237 "reserve returns the knot key the owner publishes in their did:web document" 1238 ); 1239 assert_eq!( 1240 output["key"].as_str(), 1241 Some(expected_key.as_str()), 1242 "create returns the knot-held key the owner published in their own did:web document" 1243 ); 1244 1245 assert_eq!( 1246 as_member( 1247 &world, 1248 crate::collaborators::add_collaborator, 1249 ADD_COLLAB, 1250 json!({ "repo": did, "subject": "did:web:witchcraft.systems" }) 1251 ) 1252 .await 1253 .status(), 1254 StatusCode::OK 1255 ); 1256 assert_eq!( 1257 world 1258 .state 1259 .index 1260 .is_collaborator(&repo_did, &account("witchcraft.systems")), 1261 Resolved::Ready(true), 1262 "the collaborator COB signed by the knot-held repo key lands and is visible" 1263 ); 1264 1265 let squid = "did:web:squid.olaren.dev"; 1266 let victim = RepoDid::new(squid).unwrap(); 1267 world.layout.create(&victim).unwrap(); 1268 reserve_repo_key(&world, squid).await; 1269 assert_eq!( 1270 create_status( 1271 &world, 1272 &world.member, 1273 MEMBER_HOST, 1274 json!({ "rkey": "anemone", "name": "anemone", "repoDid": squid }) 1275 ) 1276 .await, 1277 StatusCode::CONFLICT 1278 ); 1279 assert!( 1280 world.layout.open(&victim).is_ok(), 1281 "a colliding create mustn't delete the repository already on disk" 1282 ); 1283} 1284 1285#[tokio::test] 1286async fn a_rejected_plc_submission_is_a_bad_gateway() { 1287 let admin = signer(1); 1288 let key = admin.public_key().as_bytes().to_vec(); 1289 let responder = doc_responder(key, || StatusCode::BAD_REQUEST); 1290 let (_dir, state) = build_state(responder, true); 1291 let token = mint(&admin, &account(ADMIN_HOST), CREATE); 1292 let status = into_response( 1293 crate::repos::create_repo( 1294 State(Arc::clone(&state)), 1295 bearer(&token), 1296 crate::Method::from_nsid(CREATE), 1297 body(json!({ "rkey": "conch", "name": "conch" })), 1298 ) 1299 .await, 1300 ) 1301 .status(); 1302 assert_eq!( 1303 status, 1304 StatusCode::BAD_GATEWAY, 1305 "non-transient PLC rejection surfaces as 502, distinct from an internal 500" 1306 ); 1307 assert_eq!( 1308 state.secrets.len(), 1309 1, 1310 "rejected PLC submission rolls back the minted repo key, leaving only the knot's own. Unpublished did:plc orphans nothing" 1311 ); 1312 assert_eq!( 1313 state.index.resolve_repo( 1314 &OwnerDid::new(format!("did:web:{ADMIN_HOST}")).unwrap(), 1315 &RepoRkey::new("conch").unwrap() 1316 ), 1317 Resolved::Ready(None), 1318 "repo isn't registered after a rejected PLC submission" 1319 ); 1320} 1321 1322#[tokio::test] 1323async fn a_rejected_plc_submission_never_touches_the_registry() { 1324 let admin = signer(1); 1325 let key = admin.public_key().as_bytes().to_vec(); 1326 let reject_posts = Arc::new(std::sync::atomic::AtomicBool::new(false)); 1327 let reject = Arc::clone(&reject_posts); 1328 let responder = doc_responder(key, move || { 1329 if reject.load(Ordering::Relaxed) { 1330 StatusCode::BAD_REQUEST 1331 } else { 1332 StatusCode::OK 1333 } 1334 }); 1335 let (_dir, state) = build_state(responder, true); 1336 let owner = OwnerDid::new(format!("did:web:{ADMIN_HOST}")).unwrap(); 1337 1338 let mint_create = || mint(&admin, &account(ADMIN_HOST), CREATE); 1339 let anemone = || body(json!({ "rkey": "anemone", "name": "anemone" })); 1340 assert_eq!( 1341 into_response( 1342 crate::repos::create_repo( 1343 State(Arc::clone(&state)), 1344 bearer(&mint_create()), 1345 crate::Method::from_nsid(CREATE), 1346 anemone(), 1347 ) 1348 .await 1349 ) 1350 .status(), 1351 StatusCode::OK 1352 ); 1353 let victim = match state 1354 .index 1355 .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()) 1356 { 1357 Resolved::Ready(Some(did)) => did, 1358 other => panic!("victim repo wasn't registered: {other:?}"), 1359 }; 1360 1361 let token = mint(&admin, &account(ADMIN_HOST), RENAME); 1362 assert_eq!( 1363 into_response( 1364 crate::repos::rename_repo( 1365 State(Arc::clone(&state)), 1366 bearer(&token), 1367 crate::Method::from_nsid(RENAME), 1368 body(json!({ "repo": victim.as_str(), "rkey": "barnacle", "name": "barnacle" })), 1369 ) 1370 .await 1371 ) 1372 .status(), 1373 StatusCode::OK 1374 ); 1375 1376 reject_posts.store(true, Ordering::Relaxed); 1377 assert_eq!( 1378 into_response( 1379 crate::repos::create_repo( 1380 State(Arc::clone(&state)), 1381 bearer(&mint_create()), 1382 crate::Method::from_nsid(CREATE), 1383 anemone(), 1384 ) 1385 .await 1386 ) 1387 .status(), 1388 StatusCode::BAD_GATEWAY 1389 ); 1390 1391 assert_eq!( 1392 state 1393 .index 1394 .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()), 1395 Resolved::Ready(Some(victim.clone())), 1396 "stale alias still resolves to its prior holder because the failed create never registered" 1397 ); 1398 1399 state.index.refresh_registry().unwrap(); 1400 assert_eq!( 1401 state 1402 .index 1403 .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()), 1404 Resolved::Ready(Some(victim.clone())), 1405 "durable registry COB has no trace of the failed create" 1406 ); 1407 assert_eq!( 1408 state.index.rkey_of(&victim), 1409 Resolved::Ready(Some(RepoRkey::new("barnacle").unwrap())), 1410 "victim's canonical rkey is unmoved" 1411 ); 1412} 1413 1414#[tokio::test] 1415async fn resolve_by_name_matches_the_rkey_case_sensitively() { 1416 let admin = signer(1); 1417 let key = admin.public_key().as_bytes().to_vec(); 1418 let responder = doc_responder(key, || StatusCode::OK); 1419 let (_dir, state) = build_state(responder, true); 1420 let owner = OwnerDid::new(format!("did:web:{ADMIN_HOST}")).unwrap(); 1421 1422 assert_eq!( 1423 into_response( 1424 crate::repos::create_repo( 1425 State(Arc::clone(&state)), 1426 bearer(&mint(&admin, &account(ADMIN_HOST), CREATE)), 1427 crate::Method::from_nsid(CREATE), 1428 body(json!({ "rkey": "anemone", "name": "anemone" })), 1429 ) 1430 .await 1431 ) 1432 .status(), 1433 StatusCode::OK 1434 ); 1435 1436 assert!( 1437 crate::merge::resolve_by_name(&*state, &owner, "anemone").is_ok(), 1438 "the exact rkey resolves" 1439 ); 1440 assert!( 1441 crate::merge::resolve_by_name(&*state, &owner, "Anemone").is_err(), 1442 "a differently-cased name must not resolve to a distinct rkey, atproto record keys are case-sensitive" 1443 ); 1444} 1445 1446#[tokio::test] 1447async fn reserve_key_refuses_once_the_pending_limit_is_reached() { 1448 let world = World::with_pending_limit(2); 1449 add_member_helper(&world).await; 1450 1451 assert_eq!( 1452 reserve_status(&world, &world.member, MEMBER_HOST, "did:web:p0.olaren.dev").await, 1453 StatusCode::OK, 1454 "first reservation is within the limit" 1455 ); 1456 assert_eq!( 1457 reserve_status(&world, &world.member, MEMBER_HOST, "did:web:p1.olaren.dev").await, 1458 StatusCode::OK, 1459 "second reservation reaches the limit" 1460 ); 1461 assert_eq!( 1462 reserve_status(&world, &world.member, MEMBER_HOST, "did:web:p2.olaren.dev").await, 1463 StatusCode::TOO_MANY_REQUESTS, 1464 "member cannot grow the sealed store without bound past the pending-reservation limit" 1465 ); 1466} 1467 1468#[tokio::test] 1469async fn one_account_cannot_exhaust_the_global_reservation_budget() { 1470 let world = World::with_limits(256, 2); 1471 add_member_helper(&world).await; 1472 1473 let world_ref = &world; 1474 futures::stream::iter(["did:web:m0.olaren.dev", "did:web:m1.olaren.dev"]) 1475 .for_each(|did| async move { 1476 assert_eq!( 1477 reserve_status(world_ref, &world_ref.member, MEMBER_HOST, did).await, 1478 StatusCode::OK 1479 ); 1480 }) 1481 .await; 1482 assert_eq!( 1483 reserve_status(&world, &world.member, MEMBER_HOST, "did:web:m2.olaren.dev").await, 1484 StatusCode::TOO_MANY_REQUESTS, 1485 "member is held to its per-actor reservation budget" 1486 ); 1487 assert_eq!( 1488 reserve_status( 1489 &world, 1490 &world.admin, 1491 ADMIN_HOST, 1492 "did:web:admin0.olaren.dev" 1493 ) 1494 .await, 1495 StatusCode::OK, 1496 "a different account keeps its own budget while global capacity remains" 1497 ); 1498} 1499 1500#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 1501async fn concurrent_reserve_key_calls_all_succeed() { 1502 let world = World::new(); 1503 add_member_helper(&world).await; 1504 1505 let handles: Vec<_> = (0..24) 1506 .map(|i| { 1507 let state = world.state(); 1508 let token = mint(&world.member, &account(MEMBER_HOST), RESERVE); 1509 let payload = body(json!({ "repoDid": format!("did:web:r{i}.olaren.dev") })); 1510 tokio::spawn(async move { 1511 into_response( 1512 crate::repos::reserve_key( 1513 state, 1514 bearer(&token), 1515 crate::Method::from_nsid(RESERVE), 1516 payload, 1517 ) 1518 .await, 1519 ) 1520 .status() 1521 }) 1522 }) 1523 .collect(); 1524 1525 futures::future::join_all(handles) 1526 .await 1527 .into_iter() 1528 .for_each(|result| { 1529 assert_eq!( 1530 result.unwrap(), 1531 StatusCode::OK, 1532 "concurrent reserveKey mustn't race the in-memory reservation map" 1533 ) 1534 }); 1535} 1536 1537#[tokio::test] 1538async fn collaborator_lifecycle() { 1539 let world = World::new(); 1540 add_member_helper(&world).await; 1541 let repo_did = create_repo_helper(&world, "scallop").await; 1542 let subject = || json!({ "repo": repo_did.as_str(), "subject": "did:web:witchcraft.systems" }); 1543 1544 assert_eq!( 1545 as_member( 1546 &world, 1547 crate::collaborators::add_collaborator, 1548 ADD_COLLAB, 1549 subject() 1550 ) 1551 .await 1552 .status(), 1553 StatusCode::OK 1554 ); 1555 assert_eq!( 1556 world 1557 .state 1558 .index 1559 .is_collaborator(&repo_did, &account("witchcraft.systems")), 1560 Resolved::Ready(true) 1561 ); 1562 let added = last_event(&world, "sh.tangled.repo.collaboratorUpdate"); 1563 assert_eq!(added.payload["op"], "add"); 1564 assert_eq!( 1565 added.payload["subject"], 1566 account("witchcraft.systems").to_string() 1567 ); 1568 assert_eq!(added.payload["repo"], repo_did.to_string()); 1569 1570 assert_eq!( 1571 as_member( 1572 &world, 1573 crate::collaborators::remove_collaborator, 1574 REMOVE_COLLAB, 1575 subject() 1576 ) 1577 .await 1578 .status(), 1579 StatusCode::OK 1580 ); 1581 assert_eq!( 1582 world 1583 .state 1584 .index 1585 .is_collaborator(&repo_did, &account("witchcraft.systems")), 1586 Resolved::Ready(false), 1587 "removed collaborator is gone on the very next read" 1588 ); 1589 let removed = last_event(&world, "sh.tangled.repo.collaboratorUpdate"); 1590 assert_eq!(removed.payload["op"], "remove"); 1591 assert_eq!( 1592 removed.payload["subject"], 1593 account("witchcraft.systems").to_string() 1594 ); 1595 assert_eq!(removed.payload["repo"], repo_did.to_string()); 1596 1597 let baseline = event_count(&world); 1598 assert_eq!( 1599 as_member( 1600 &world, 1601 crate::collaborators::remove_collaborator, 1602 REMOVE_COLLAB, 1603 subject() 1604 ) 1605 .await 1606 .status(), 1607 StatusCode::OK 1608 ); 1609 assert_eq!( 1610 event_count(&world), 1611 baseline, 1612 "removing a non-collaborator is a no-op and emits no event" 1613 ); 1614} 1615 1616#[tokio::test] 1617async fn repo_management_is_owner_or_collaborator_gated() { 1618 let world = World::new(); 1619 add_member_helper(&world).await; 1620 let repo_did = create_repo_helper(&world, "squid").await; 1621 let at = format!("at://did:web:{MEMBER_HOST}/sh.tangled.repo/squid"); 1622 let did = format!("did:web:{MEMBER_HOST}"); 1623 1624 assert_eq!( 1625 as_admin( 1626 &world, 1627 crate::collaborators::add_collaborator, 1628 ADD_COLLAB, 1629 json!({ "repo": repo_did.as_str(), "subject": "did:web:isabelroses.com" }) 1630 ) 1631 .await 1632 .status(), 1633 StatusCode::FORBIDDEN, 1634 "collaborator management is the repo owner's right instead of a knot admin's" 1635 ); 1636 assert_eq!( 1637 as_stranger( 1638 &world, 1639 crate::branches::set_default_branch, 1640 SET_DEFAULT, 1641 json!({ "repo": at, "defaultBranch": "trunk" }) 1642 ) 1643 .await 1644 .status(), 1645 StatusCode::FORBIDDEN, 1646 "a stranger cannot set the default branch" 1647 ); 1648 assert_eq!( 1649 as_stranger( 1650 &world, 1651 crate::branches::delete_branch, 1652 DELETE_BRANCH, 1653 json!({ "repo": at, "branch": "trunk" }) 1654 ) 1655 .await 1656 .status(), 1657 StatusCode::FORBIDDEN, 1658 "a stranger cannot delete a branch" 1659 ); 1660 assert_eq!( 1661 as_stranger( 1662 &world, 1663 crate::repos::delete_repo, 1664 DELETE, 1665 json!({ "did": did, "name": "squid", "rkey": "squid" }) 1666 ) 1667 .await 1668 .status(), 1669 StatusCode::FORBIDDEN, 1670 "neither owner nor a knot admin, so delete is refused" 1671 ); 1672 1673 assert_eq!( 1674 rename_repo_as(&world, &world.stranger, STRANGER_HOST, &repo_did, "stolen").await, 1675 StatusCode::FORBIDDEN 1676 ); 1677 assert_eq!( 1678 world.state.index.rkey_of(&repo_did), 1679 Resolved::Ready(Some(RepoRkey::new("squid").unwrap())), 1680 "the canonical rkey is untouched by the rejected rename" 1681 ); 1682 1683 assert_eq!( 1684 as_member( 1685 &world, 1686 crate::collaborators::add_collaborator, 1687 ADD_COLLAB, 1688 json!({ "repo": repo_did.as_str(), "subject": format!("did:web:{STRANGER_HOST}") }) 1689 ) 1690 .await 1691 .status(), 1692 StatusCode::OK 1693 ); 1694 assert_eq!( 1695 rename_repo_as( 1696 &world, 1697 &world.stranger, 1698 STRANGER_HOST, 1699 &repo_did, 1700 "periwinkle" 1701 ) 1702 .await, 1703 StatusCode::OK, 1704 "rename is gated by can_push, so a collaborator may rename" 1705 ); 1706} 1707 1708#[tokio::test] 1709async fn the_owner_sets_the_default_branch_through_an_at_uri() { 1710 let world = World::new(); 1711 add_member_helper(&world).await; 1712 let repo_did = create_repo_helper(&world, "mussel").await; 1713 1714 assert_eq!( 1715 as_member(&world, crate::branches::set_default_branch, SET_DEFAULT, json!({ "repo": format!("at://did:web:{MEMBER_HOST}/sh.tangled.repo/mussel"), "defaultBranch": "trunk" })).await.status(), 1716 StatusCode::OK 1717 ); 1718 1719 let repo = world.layout.open(&repo_did).unwrap(); 1720 assert_eq!( 1721 repo.default_branch().map(|name| name.as_str().to_string()), 1722 Some("refs/heads/trunk".to_string()), 1723 "HEAD now points at the requested default branch" 1724 ); 1725 1726 let event = only_git_event(&world); 1727 assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); 1728 assert_eq!(event.payload["repo"], repo_did.to_string()); 1729 assert_eq!(event.payload["ownerDid"], account(MEMBER_HOST).to_string()); 1730 assert_eq!( 1731 event.payload["committerDid"], 1732 account(MEMBER_HOST).to_string(), 1733 "the actor who set the default branch is the committer on the wire" 1734 ); 1735} 1736 1737#[tokio::test] 1738async fn set_default_branch_rejections() { 1739 use knot_cob::CobStore; 1740 use knot_cobs::{CollaboratorsChange, Grant}; 1741 use knot_git::RefUpdate; 1742 use knot_types::RefName; 1743 1744 let world = World::new(); 1745 add_member_helper(&world).await; 1746 let repo_did = create_repo_helper(&world, "mussel").await; 1747 1748 let git = world.layout.open(&repo_did).unwrap(); 1749 let now = world.state.now(); 1750 let signer = world.state.secrets.signer(&world.state.knot_did).unwrap(); 1751 let created = CobStore::new(&git) 1752 .create( 1753 &knot_cob::CobHome::from(&repo_did), 1754 &CollaboratorsChange::Add(Grant { 1755 subject: account("olaren.dev"), 1756 added_by: account("olaren.dev"), 1757 created_at: now, 1758 }), 1759 &signer, 1760 now, 1761 ) 1762 .unwrap(); 1763 git.update_ref(&RefUpdate::Create { 1764 name: RefName::new("refs/heads/main").unwrap(), 1765 new: created.tip.oid(), 1766 }) 1767 .unwrap(); 1768 1769 assert_eq!( 1770 as_member(&world, crate::branches::set_default_branch, SET_DEFAULT, json!({ "repo": format!("at://did:web:{MEMBER_HOST}/sh.tangled.notrepo/mussel"), "defaultBranch": "trunk" })).await.status(), 1771 StatusCode::BAD_REQUEST, 1772 "an at-uri addressing a collection other than sh.tangled.repo is rejected" 1773 ); 1774 assert_eq!( 1775 as_member(&world, crate::branches::set_default_branch, SET_DEFAULT, json!({ "repo": format!("at://did:web:{MEMBER_HOST}/sh.tangled.repo/mussel"), "defaultBranch": "ghost" })).await.status(), 1776 StatusCode::NOT_FOUND, 1777 "a populated repo rejects a default pointing at a branch that doesn't exist" 1778 ); 1779} 1780 1781#[tokio::test] 1782async fn delete_branch_removes_a_branch_and_refuses_the_default() { 1783 use knot_cob::CobStore; 1784 use knot_cobs::{CollaboratorsChange, Grant}; 1785 use knot_git::RefUpdate; 1786 use knot_types::RefName; 1787 1788 let world = World::new(); 1789 add_member_helper(&world).await; 1790 let repo_did = create_repo_helper(&world, "periwinkle").await; 1791 let git = world.layout.open(&repo_did).unwrap(); 1792 let now = world.state.now(); 1793 let signer = world.state.secrets.signer(&world.state.knot_did).unwrap(); 1794 let created = CobStore::new(&git) 1795 .create( 1796 &knot_cob::CobHome::from(&repo_did), 1797 &CollaboratorsChange::Add(Grant { 1798 subject: account("olaren.dev"), 1799 added_by: account("olaren.dev"), 1800 created_at: now, 1801 }), 1802 &signer, 1803 now, 1804 ) 1805 .unwrap(); 1806 let oid = created.tip.oid(); 1807 ["refs/heads/main", "refs/heads/trunk"] 1808 .into_iter() 1809 .for_each(|name| { 1810 git.update_ref(&RefUpdate::Create { 1811 name: RefName::new(name).unwrap(), 1812 new: oid, 1813 }) 1814 .unwrap(); 1815 }); 1816 git.set_head(&RefName::new("refs/heads/main").unwrap()) 1817 .unwrap(); 1818 1819 let at = format!("at://did:web:{MEMBER_HOST}/sh.tangled.repo/periwinkle"); 1820 assert_eq!( 1821 as_member( 1822 &world, 1823 crate::branches::delete_branch, 1824 DELETE_BRANCH, 1825 json!({ "repo": at, "branch": "trunk" }) 1826 ) 1827 .await 1828 .status(), 1829 StatusCode::OK, 1830 "a non-default branch is deleted" 1831 ); 1832 assert!( 1833 git.find_ref(&RefName::new("refs/heads/trunk").unwrap()) 1834 .unwrap() 1835 .is_none(), 1836 "trunk is gone" 1837 ); 1838 assert_eq!( 1839 as_member( 1840 &world, 1841 crate::branches::delete_branch, 1842 DELETE_BRANCH, 1843 json!({ "repo": at, "branch": "main" }) 1844 ) 1845 .await 1846 .status(), 1847 StatusCode::BAD_REQUEST, 1848 "the current default branch cannot be deleted" 1849 ); 1850 1851 let event = only_git_event(&world); 1852 assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); 1853 assert_eq!(event.payload["repo"], repo_did.to_string()); 1854 assert_eq!(event.payload["ref"], "refs/heads/trunk"); 1855 assert_eq!( 1856 event.payload["oldSha"], 1857 oid.to_string(), 1858 "the deletion event includes the branch's old tip" 1859 ); 1860 assert_eq!( 1861 event.payload["newSha"], 1862 git.object_format().null_oid().to_string(), 1863 "deletion reports the null oid as the new sha" 1864 ); 1865 assert_eq!( 1866 event.payload["committerDid"], 1867 account(MEMBER_HOST).to_string() 1868 ); 1869} 1870 1871#[tokio::test] 1872async fn delete_repo_lifecycle_and_guards() { 1873 let world = World::new(); 1874 add_member_helper(&world).await; 1875 let did = format!("did:web:{MEMBER_HOST}"); 1876 1877 let plain = create_repo_helper(&world, "whelk").await; 1878 assert_eq!( 1879 as_member( 1880 &world, 1881 crate::repos::delete_repo, 1882 DELETE, 1883 json!({ "did": did, "name": "whelk", "rkey": "whelk" }) 1884 ) 1885 .await 1886 .status(), 1887 StatusCode::OK 1888 ); 1889 assert!( 1890 world.layout.open(&plain).is_err(), 1891 "bare repo is removed from disk" 1892 ); 1893 assert_eq!( 1894 resolve(&world, "whelk"), 1895 Resolved::Ready(None), 1896 "repo is deregistered" 1897 ); 1898 1899 let guarded = create_repo_helper(&world, "conch").await; 1900 world.publish_pds_record("conch"); 1901 let delete_conch = || json!({ "did": did, "name": "conch", "rkey": "conch" }); 1902 assert_eq!( 1903 as_member(&world, crate::repos::delete_repo, DELETE, delete_conch()) 1904 .await 1905 .status(), 1906 StatusCode::CONFLICT, 1907 "the guard refuses while the sh.tangled.repo record is still on the owner's PDS" 1908 ); 1909 assert!( 1910 world.layout.open(&guarded).is_ok(), 1911 "a refused delete left the repo intact on disk" 1912 ); 1913 1914 let force_conch = || json!({ "did": did, "name": "conch", "rkey": "conch", "force": true }); 1915 assert_eq!( 1916 as_member(&world, crate::repos::delete_repo, DELETE, force_conch()) 1917 .await 1918 .status(), 1919 StatusCode::FORBIDDEN, 1920 "force is an admin-only escape hatch instead of the owner's" 1921 ); 1922 assert_eq!( 1923 as_admin(&world, crate::repos::delete_repo, DELETE, force_conch()) 1924 .await 1925 .status(), 1926 StatusCode::OK, 1927 "a knot admin forces the delete past the lingering PDS record" 1928 ); 1929 assert!( 1930 world.layout.open(&guarded).is_err(), 1931 "forced delete removed the repo from disk" 1932 ); 1933} 1934 1935#[tokio::test] 1936async fn rename_alias_lifecycle() { 1937 let world = World::new(); 1938 add_member_helper(&world).await; 1939 1940 let repo_a = create_repo_helper(&world, "alpha").await; 1941 assert_eq!( 1942 rename_repo_as(&world, &world.member, MEMBER_HOST, &repo_a, "alphanew").await, 1943 StatusCode::OK 1944 ); 1945 assert_eq!( 1946 resolve(&world, "alphanew"), 1947 Resolved::Ready(Some(repo_a.clone())), 1948 "the new rkey resolves on the very next read" 1949 ); 1950 assert_eq!( 1951 resolve(&world, "alpha"), 1952 Resolved::Ready(Some(repo_a.clone())), 1953 "the prior rkey keeps resolving as an alias" 1954 ); 1955 assert_eq!( 1956 world.state.index.rkey_of(&repo_a), 1957 Resolved::Ready(Some(RepoRkey::new("alphanew").unwrap())), 1958 "the new rkey is canonical" 1959 ); 1960 1961 assert_eq!( 1962 as_member(&world, crate::branches::set_default_branch, SET_DEFAULT, json!({ "repo": format!("at://did:web:{MEMBER_HOST}/sh.tangled.repo/alpha"), "defaultBranch": "trunk" })).await.status(), 1963 StatusCode::OK, 1964 "an at-uri with the pre-rename rkey still reaches the repo" 1965 ); 1966 1967 let repo_a2 = create_repo_helper(&world, "alpha").await; 1968 assert_ne!(repo_a2, repo_a, "a brand-new repo DID was minted"); 1969 assert_eq!( 1970 resolve(&world, "alpha"), 1971 Resolved::Ready(Some(repo_a2.clone())), 1972 "the reused rkey now resolves to the new repo" 1973 ); 1974 assert_eq!( 1975 resolve(&world, "alphanew"), 1976 Resolved::Ready(Some(repo_a.clone())), 1977 "the renamed repo keeps its canonical rkey" 1978 ); 1979 1980 let repo_b = create_repo_helper(&world, "beta").await; 1981 assert_eq!( 1982 rename_repo_as(&world, &world.member, MEMBER_HOST, &repo_b, "alpha").await, 1983 StatusCode::CONFLICT, 1984 "a rename cannot take the canonical rkey of another live repo" 1985 ); 1986 assert_eq!( 1987 resolve(&world, "alpha"), 1988 Resolved::Ready(Some(repo_a2)), 1989 "the contested rkey still belongs to its original repo" 1990 ); 1991 1992 let repo_g = create_repo_helper(&world, "gamma").await; 1993 assert_eq!( 1994 rename_repo_as(&world, &world.member, MEMBER_HOST, &repo_g, "gammanew").await, 1995 StatusCode::OK 1996 ); 1997 assert_eq!( 1998 as_member(&world, crate::repos::delete_repo, DELETE, json!({ "did": format!("did:web:{MEMBER_HOST}"), "name": "gammanew", "rkey": "gammanew" })).await.status(), 1999 StatusCode::OK 2000 ); 2001 assert!( 2002 world.layout.open(&repo_g).is_err(), 2003 "a renamed repo is deleted through its new rkey" 2004 ); 2005 assert_eq!( 2006 resolve(&world, "gamma"), 2007 Resolved::Ready(None), 2008 "deletion drops the retained alias along with the repo" 2009 ); 2010 2011 let ghost = RepoDid::new("did:web:ghost.nel.pet").unwrap(); 2012 assert_eq!( 2013 rename_repo_as(&world, &world.member, MEMBER_HOST, &ghost, "kelp").await, 2014 StatusCode::NOT_FOUND, 2015 "renaming a repo this knot doesn't host is 404 instead of 403" 2016 ); 2017} 2018 2019#[tokio::test] 2020async fn a_rename_against_a_warming_registry_is_unavailable_not_forbidden() { 2021 let admin = signer(1); 2022 let responder = doc_responder(admin.public_key().as_bytes().to_vec(), || StatusCode::OK); 2023 let (_dir, state) = build_state(responder, false); 2024 2025 let token = mint(&admin, &account(ADMIN_HOST), RENAME); 2026 assert_eq!( 2027 into_response( 2028 crate::repos::rename_repo( 2029 State(Arc::clone(&state)), 2030 bearer(&token), 2031 crate::Method::from_nsid(RENAME), 2032 body(json!({ "repo": "did:web:squid.nel.pet", "rkey": "kelp", "name": "kelp" })), 2033 ) 2034 .await 2035 ) 2036 .status(), 2037 StatusCode::SERVICE_UNAVAILABLE, 2038 "a warming registry is a retryable 503, never a permanent 403" 2039 ); 2040} 2041 2042#[tokio::test] 2043async fn the_router_sheds_a_pre_auth_flood_from_one_peer() { 2044 use axum::body::Body; 2045 use axum::extract::ConnectInfo; 2046 use std::net::SocketAddr; 2047 use tower::ServiceExt; 2048 2049 let world = World::new(); 2050 let app = crate::router(Arc::clone(&world.state)); 2051 let peer = SocketAddr::from(([203, 0, 113, 7], 5555)); 2052 2053 let statuses: Vec<StatusCode> = futures::stream::iter(0..22) 2054 .then(|_| { 2055 let app = app.clone(); 2056 async move { 2057 let mut request = http::Request::builder() 2058 .method("POST") 2059 .uri(crate::members::ADD_ROUTE) 2060 .body(Body::empty()) 2061 .unwrap(); 2062 request.extensions_mut().insert(ConnectInfo(peer)); 2063 app.oneshot(request).await.unwrap().status() 2064 } 2065 }) 2066 .collect() 2067 .await; 2068 2069 assert!( 2070 statuses[..20] 2071 .iter() 2072 .all(|status| *status == StatusCode::UNAUTHORIZED), 2073 "per-peer burst is admitted and then fails auth on the missing token, got {statuses:?}" 2074 ); 2075 assert!( 2076 statuses[20..] 2077 .iter() 2078 .all(|status| *status == StatusCode::TOO_MANY_REQUESTS), 2079 "past the burst the router sheds the flood before it can reach the resolver, got {statuses:?}" 2080 ); 2081} 2082 2083#[tokio::test] 2084async fn the_router_binds_each_route_to_its_matched_method_scope() { 2085 use axum::body::Body; 2086 use axum::extract::ConnectInfo; 2087 use std::net::SocketAddr; 2088 use tower::ServiceExt; 2089 2090 let world = World::new(); 2091 let app = crate::router(Arc::clone(&world.state)); 2092 let peer = SocketAddr::from(([203, 0, 113, 23], 5555)); 2093 2094 let post = |token: String| { 2095 let mut request = http::Request::builder() 2096 .method("POST") 2097 .uri(crate::members::ADD_ROUTE) 2098 .header(http::header::AUTHORIZATION, format!("Bearer {token}")) 2099 .body(Body::from( 2100 serde_json::to_vec(&json!({ "subject": format!("did:web:{MEMBER_HOST}") })) 2101 .unwrap(), 2102 )) 2103 .unwrap(); 2104 request.extensions_mut().insert(ConnectInfo(peer)); 2105 request 2106 }; 2107 2108 let matched = mint(&world.admin, &account(ADMIN_HOST), ADD_MEMBER); 2109 assert_eq!( 2110 app.clone().oneshot(post(matched)).await.unwrap().status(), 2111 StatusCode::OK, 2112 "a token whose lxm is the route's own method authenticates" 2113 ); 2114 2115 let sibling = mint(&world.admin, &account(ADMIN_HOST), REMOVE_MEMBER); 2116 assert_eq!( 2117 app.oneshot(post(sibling)).await.unwrap().status(), 2118 StatusCode::UNAUTHORIZED, 2119 "a token minted for a sibling method is rejected at the addMember route" 2120 ); 2121} 2122 2123#[tokio::test] 2124async fn the_http_push_surface_sheds_a_bogus_credential_flood_from_one_peer() { 2125 use axum::body::Body; 2126 use axum::extract::ConnectInfo; 2127 use base64::Engine as _; 2128 use std::net::SocketAddr; 2129 use tower::ServiceExt; 2130 2131 let world = World::new(); 2132 add_member_helper(&world).await; 2133 let repo = create_repo_helper(&world, "kelp").await; 2134 let app = crate::router(Arc::clone(&world.state)); 2135 let peer = SocketAddr::from(([203, 0, 113, 11], 5555)); 2136 let credential = format!( 2137 "Basic {}", 2138 base64::engine::general_purpose::STANDARD.encode("git:not-a-service-jwt") 2139 ); 2140 2141 let statuses: Vec<StatusCode> = futures::stream::iter(0..22) 2142 .then(|_| { 2143 let app = app.clone(); 2144 let uri = format!("/{}/git-receive-pack", repo.as_str()); 2145 let credential = credential.clone(); 2146 async move { 2147 let mut request = http::Request::builder() 2148 .method("POST") 2149 .uri(uri) 2150 .header(http::header::AUTHORIZATION, credential) 2151 .body(Body::empty()) 2152 .unwrap(); 2153 request.extensions_mut().insert(ConnectInfo(peer)); 2154 app.oneshot(request).await.unwrap().status() 2155 } 2156 }) 2157 .collect() 2158 .await; 2159 2160 assert!( 2161 statuses[..20] 2162 .iter() 2163 .all(|status| *status == StatusCode::UNAUTHORIZED), 2164 "bogus credentials inside the burst fail authentication, got {statuses:?}" 2165 ); 2166 assert!( 2167 statuses[20..] 2168 .iter() 2169 .all(|status| *status == StatusCode::TOO_MANY_REQUESTS), 2170 "past the burst the push surface sheds the flood before it can reach the resolver, got {statuses:?}" 2171 ); 2172} 2173 2174#[tokio::test] 2175async fn health_is_unauthenticated_and_exempt_from_shedding() { 2176 use axum::body::Body; 2177 use axum::extract::ConnectInfo; 2178 use std::net::SocketAddr; 2179 use tower::ServiceExt; 2180 2181 let world = World::new(); 2182 let app = crate::router(Arc::clone(&world.state)); 2183 let peer = SocketAddr::from(([203, 0, 113, 9], 5555)); 2184 let health = || { 2185 let mut request = http::Request::builder() 2186 .method("GET") 2187 .uri(crate::service::HEALTH_ROUTE) 2188 .body(Body::empty()) 2189 .unwrap(); 2190 request.extensions_mut().insert(ConnectInfo(peer)); 2191 request 2192 }; 2193 2194 let statuses = futures::future::join_all((0..30).map(|_| { 2195 let app = app.clone(); 2196 async move { app.oneshot(health()).await.unwrap() } 2197 })) 2198 .await; 2199 assert!( 2200 statuses 2201 .iter() 2202 .all(|response| response.status() == StatusCode::OK), 2203 "health stays 200 even past the pre-auth burst, got {:?}", 2204 statuses.iter().map(|r| r.status()).collect::<Vec<_>>() 2205 ); 2206 2207 let wire = json_of(app.oneshot(health()).await.unwrap()).await; 2208 assert!( 2209 wire["version"] 2210 .as_str() 2211 .is_some_and(|v| v.starts_with("knot ")), 2212 "health reports a knot version, got {wire}" 2213 ); 2214} 2215 2216mod merge_endpoints { 2217 use super::*; 2218 use std::path::Path; 2219 2220 use knot_git::{EntryKind, Identity, NewCommit, RefUpdate, StagedAction, StagedChange}; 2221 use knot_types::{Oid, RefName, UnixSeconds}; 2222 2223 const EMPTY_TREE: &str = "4b825dc642cb6eb9a060e54bf8d69288fbee4904"; 2224 const MERGE: &str = "sh.tangled.repo.merge"; 2225 2226 const UNIFIED_PATCH: &str = concat!( 2227 "diff --git a/reef.txt b/reef.txt\n", 2228 "index 1111111..2222222 100644\n", 2229 "--- a/reef.txt\n", 2230 "+++ b/reef.txt\n", 2231 "@@ -1 +1 @@\n", 2232 "-old line\n", 2233 "+new line\n", 2234 ); 2235 2236 const CONFLICTING_PATCH: &str = concat!( 2237 "diff --git a/reef.txt b/reef.txt\n", 2238 "index 1111111..2222222 100644\n", 2239 "--- a/reef.txt\n", 2240 "+++ b/reef.txt\n", 2241 "@@ -1 +1 @@\n", 2242 "-something else entirely\n", 2243 "+new line\n", 2244 ); 2245 2246 fn seed_main(world: &World, repo_did: &RepoDid, files: &[(&str, &str)]) -> Oid { 2247 let repo = world.layout.open(repo_did).unwrap(); 2248 let staged: Vec<StagedChange> = files 2249 .iter() 2250 .map(|(path, content)| StagedChange { 2251 path: knot_types::RepoPath::new(*path).unwrap(), 2252 action: StagedAction::Put { 2253 content: content.as_bytes().to_vec(), 2254 kind: EntryKind::Blob, 2255 }, 2256 }) 2257 .collect(); 2258 let tree = repo 2259 .write_staged_tree(Oid::from_hex(EMPTY_TREE).unwrap(), &staged) 2260 .unwrap(); 2261 let nel = Identity { 2262 name: AuthorName::new("nel"), 2263 email: Email::new("nel@oyster.cafe"), 2264 time: UnixSeconds::new(1_000), 2265 offset_seconds: 0, 2266 }; 2267 let commit = repo 2268 .write_commit(&NewCommit { 2269 tree, 2270 parents: Vec::new(), 2271 author: nel.clone(), 2272 committer: nel, 2273 message: "base".to_string(), 2274 extra_headers: Vec::new(), 2275 }) 2276 .unwrap(); 2277 repo.update_ref(&RefUpdate::Create { 2278 name: RefName::new("refs/heads/main").unwrap(), 2279 new: commit, 2280 }) 2281 .unwrap(); 2282 commit 2283 } 2284 2285 fn main_tip(world: &World, repo_did: &RepoDid) -> Oid { 2286 world 2287 .layout 2288 .open(repo_did) 2289 .unwrap() 2290 .find_ref(&RefName::new("refs/heads/main").unwrap()) 2291 .unwrap() 2292 .unwrap() 2293 } 2294 2295 fn file_count(dir: &Path) -> usize { 2296 std::fs::read_dir(dir) 2297 .map(|entries| { 2298 entries 2299 .flatten() 2300 .map(|entry| match entry.file_type() { 2301 Ok(kind) if kind.is_dir() => file_count(&entry.path()), 2302 _ => 1, 2303 }) 2304 .sum() 2305 }) 2306 .unwrap_or(0) 2307 } 2308 2309 fn blob_at(repo: &knot_git::Repo, commit: Oid, path: &str) -> Vec<u8> { 2310 let entry = repo 2311 .entry_at(commit, &knot_types::RepoPath::new(path).unwrap()) 2312 .unwrap() 2313 .unwrap(); 2314 repo.read_blob(entry.oid).unwrap() 2315 } 2316 2317 #[tokio::test] 2318 async fn the_owner_merges_a_unified_patch_natively_and_cleans_up() { 2319 let world = World::new(); 2320 add_member_helper(&world).await; 2321 let repo_did = create_repo_helper(&world, "kelp").await; 2322 let base = seed_main(&world, &repo_did, &[("reef.txt", "old line\n")]); 2323 2324 assert_eq!( 2325 as_member( 2326 &world, 2327 crate::merge::merge, 2328 MERGE, 2329 json!({ 2330 "did": format!("did:web:{MEMBER_HOST}"), 2331 "name": "kelp", 2332 "branch": "main", 2333 "patch": UNIFIED_PATCH, 2334 "commitMessage": "Merge tide", 2335 "commitBody": "body text", 2336 "authorName": "bailey", 2337 "authorEmail": "bailey@nel.pet", 2338 }) 2339 ) 2340 .await 2341 .status(), 2342 StatusCode::OK 2343 ); 2344 2345 let repo = world.layout.open(&repo_did).unwrap(); 2346 let tip = main_tip(&world, &repo_did); 2347 assert_ne!(tip, base); 2348 let commit = repo.find_commit(tip).unwrap(); 2349 assert_eq!(commit.parents, vec![base]); 2350 assert_eq!(commit.author.name.as_str(), "bailey"); 2351 assert_eq!(commit.author.email.as_str(), "bailey@nel.pet"); 2352 assert_eq!(commit.committer.name.as_str(), "Tangled"); 2353 assert_eq!(commit.committer.email.as_str(), "noreply@tangled.sh"); 2354 assert_eq!(commit.message, "Merge tide\n\nbody text\n"); 2355 assert_eq!(blob_at(&repo, tip, "reef.txt"), b"new line\n"); 2356 2357 let event = only_git_event(&world); 2358 assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); 2359 assert_eq!(event.payload["repo"], repo_did.to_string()); 2360 assert_eq!(event.payload["ref"], "refs/heads/main"); 2361 assert_eq!(event.payload["oldSha"], base.to_string()); 2362 assert_eq!(event.payload["newSha"], tip.to_string()); 2363 assert_eq!( 2364 event.payload["committerDid"], 2365 account(MEMBER_HOST).to_string() 2366 ); 2367 2368 let staging = std::fs::read_dir(repo.path()) 2369 .unwrap() 2370 .flatten() 2371 .filter(|entry| { 2372 entry 2373 .file_name() 2374 .to_str() 2375 .is_some_and(|name| name.starts_with(knot_git::INCOMING_PREFIX)) 2376 }) 2377 .count(); 2378 assert_eq!( 2379 staging, 0, 2380 "a completed merge cleans up its staging directory" 2381 ); 2382 } 2383 2384 #[tokio::test] 2385 async fn a_native_merge_advances_the_branch_without_a_pipeline_event() { 2386 let world = World::new(); 2387 add_member_helper(&world).await; 2388 let repo_did = create_repo_helper(&world, "kelp").await; 2389 seed_main( 2390 &world, 2391 &repo_did, 2392 &[ 2393 ("reef.txt", "old line\n"), 2394 ( 2395 ".tangled/workflows/ci.yml", 2396 "engine: nixery.dev/x\nwhen:\n - event: push\n branch: ['**']\n", 2397 ), 2398 ], 2399 ); 2400 2401 assert_eq!( 2402 as_member( 2403 &world, 2404 crate::merge::merge, 2405 MERGE, 2406 json!({ 2407 "did": format!("did:web:{MEMBER_HOST}"), 2408 "name": "kelp", 2409 "branch": "main", 2410 "patch": UNIFIED_PATCH, 2411 "commitMessage": "Merge tide", 2412 "authorName": "bailey", 2413 "authorEmail": "bailey@nel.pet", 2414 }) 2415 ) 2416 .await 2417 .status(), 2418 StatusCode::OK 2419 ); 2420 2421 let update = last_event(&world, "sh.tangled.git.refUpdate"); 2422 assert_eq!(update.payload["ref"], "refs/heads/main"); 2423 assert!( 2424 replay(&world) 2425 .iter() 2426 .all(|event| event.nsid != "sh.tangled.pipeline"), 2427 "the knot emits no pipeline record" 2428 ); 2429 } 2430 2431 #[tokio::test] 2432 async fn a_format_patch_merge_creates_one_commit_per_patch_with_change_id() { 2433 let world = World::new(); 2434 add_member_helper(&world).await; 2435 let repo_did = create_repo_helper(&world, "limpet").await; 2436 let base = seed_main(&world, &repo_did, &[("reef.txt", "one\n")]); 2437 2438 let mbox = concat!( 2439 "From 1111111111111111111111111111111111111111 Mon Sep 17 00:00:00 2001\n", 2440 "From: olaren <olaren@olaren.dev>\n", 2441 "Date: Tue, 5 Sep 2023 12:00:00 +0530\n", 2442 "Subject: [PATCH 1/2] first step\n", 2443 "\n", 2444 "step one body\n", 2445 "---\n", 2446 " reef.txt | 2 +-\n", 2447 "\n", 2448 "diff --git a/reef.txt b/reef.txt\n", 2449 "index 1111111..2222222 100644\n", 2450 "--- a/reef.txt\n", 2451 "+++ b/reef.txt\n", 2452 "@@ -1 +1 @@\n", 2453 "-one\n", 2454 "+two\n", 2455 "-- \n2.43.0\n\n", 2456 "From 2222222222222222222222222222222222222222 Mon Sep 17 00:00:00 2001\n", 2457 "From: olaren <olaren@olaren.dev>\n", 2458 "Date: Tue, 5 Sep 2023 13:00:00 +0530\n", 2459 "Subject: [PATCH 2/2] second step\n", 2460 "Change-Id: Ifeedfacecafe\n", 2461 "\n", 2462 "---\n", 2463 "diff --git a/reef.txt b/reef.txt\n", 2464 "index 2222222..3333333 100644\n", 2465 "--- a/reef.txt\n", 2466 "+++ b/reef.txt\n", 2467 "@@ -1 +1 @@\n", 2468 "-two\n", 2469 "+three\n", 2470 ); 2471 2472 assert_eq!( 2473 as_member( 2474 &world, 2475 crate::merge::merge, 2476 MERGE, 2477 json!({ 2478 "did": format!("did:web:{MEMBER_HOST}"), 2479 "name": "limpet", 2480 "branch": "main", 2481 "patch": mbox, 2482 }) 2483 ) 2484 .await 2485 .status(), 2486 StatusCode::OK 2487 ); 2488 2489 let repo = world.layout.open(&repo_did).unwrap(); 2490 let tip = main_tip(&world, &repo_did); 2491 let second = repo.find_commit(tip).unwrap(); 2492 assert_eq!(second.message, "second step\n"); 2493 assert_eq!( 2494 second.change_id(), 2495 Some(knot_git::CommitChangeId::new("Ifeedfacecafe").unwrap()) 2496 ); 2497 assert_eq!(second.author.name.as_str(), "olaren"); 2498 assert_eq!(second.author.email.as_str(), "olaren@olaren.dev"); 2499 assert_eq!(second.author.time.get(), 1_693_899_000); 2500 assert_eq!(second.author.offset_seconds, 19_800); 2501 assert_eq!(second.committer.name.as_str(), "Tangled"); 2502 assert_eq!(second.committer.email.as_str(), "noreply@tangled.sh"); 2503 2504 let first = repo.find_commit(second.parents[0]).unwrap(); 2505 assert_eq!(first.message, "first step\n\nstep one body\n"); 2506 assert_eq!(first.author.time.get(), 1_693_895_400); 2507 assert_eq!(first.parents, vec![base]); 2508 assert_eq!(blob_at(&repo, tip, "reef.txt"), b"three\n"); 2509 } 2510 2511 #[tokio::test] 2512 async fn merge_rejections_move_nothing() { 2513 let world = World::new(); 2514 add_member_helper(&world).await; 2515 let repo_did = create_repo_helper(&world, "scallop").await; 2516 let base = seed_main(&world, &repo_did, &[("reef.txt", "old line\n")]); 2517 let did = format!("did:web:{MEMBER_HOST}"); 2518 2519 assert_eq!( 2520 as_stranger( 2521 &world, 2522 crate::merge::merge, 2523 MERGE, 2524 json!({ "did": did, "name": "scallop", "branch": "main", "patch": UNIFIED_PATCH }) 2525 ) 2526 .await 2527 .status(), 2528 StatusCode::FORBIDDEN, 2529 "a stranger cannot merge" 2530 ); 2531 2532 assert_eq!( 2533 as_member( 2534 &world, 2535 crate::merge::merge, 2536 MERGE, 2537 json!({ "did": did, "name": "scallop", "branch": "main", "patch": UNIFIED_PATCH }) 2538 ) 2539 .await 2540 .status(), 2541 StatusCode::BAD_REQUEST, 2542 "a merge without a commit message is rejected" 2543 ); 2544 assert_eq!( 2545 main_tip(&world, &repo_did), 2546 base, 2547 "rejected merge mustn't move the branch" 2548 ); 2549 2550 assert_eq!( 2551 as_member(&world, crate::merge::merge, MERGE, json!({ "did": did, "name": "scallop", "branch": "driftwood", "patch": UNIFIED_PATCH, "commitMessage": "tide" })).await.status(), 2552 StatusCode::BAD_REQUEST, 2553 "merging into a branch the repo lacks is an invalid request" 2554 ); 2555 2556 let response = as_member(&world, crate::merge::merge, MERGE, json!({ "did": did, "name": "scallop", "branch": "main", "patch": CONFLICTING_PATCH, "commitMessage": "tide" })).await; 2557 assert_eq!(response.status(), StatusCode::CONFLICT); 2558 let json = json_of(response).await; 2559 assert_eq!(json["error"], "MergeConflict"); 2560 assert!( 2561 json["message"] 2562 .as_str() 2563 .unwrap() 2564 .starts_with("Merge failed due to conflicts"), 2565 ); 2566 assert_eq!( 2567 main_tip(&world, &repo_did), 2568 base, 2569 "a conflicted merge mustn't move the branch" 2570 ); 2571 } 2572 2573 #[tokio::test] 2574 async fn merge_check_is_open_and_reports_clean_conflicted_and_broken() { 2575 let world = World::new(); 2576 add_member_helper(&world).await; 2577 let repo_did = create_repo_helper(&world, "scallop").await; 2578 let base = seed_main(&world, &repo_did, &[("reef.txt", "old line\n")]); 2579 2580 let repo = world.layout.open(&repo_did).unwrap(); 2581 let objects_before = file_count(&repo.objects_dir()); 2582 2583 let input = |patch: &str| { 2584 body(json!({ 2585 "did": format!("did:web:{MEMBER_HOST}"), 2586 "name": "scallop", 2587 "branch": "main", 2588 "patch": patch, 2589 })) 2590 }; 2591 2592 let clean = json_of( 2593 crate::merge::merge_check(world.state(), input(UNIFIED_PATCH)) 2594 .await 2595 .unwrap(), 2596 ) 2597 .await; 2598 assert_eq!(clean["is_conflicted"], false); 2599 assert!(clean.get("conflicts").is_none()); 2600 2601 let conflicted = json_of( 2602 crate::merge::merge_check(world.state(), input(CONFLICTING_PATCH)) 2603 .await 2604 .unwrap(), 2605 ) 2606 .await; 2607 assert_eq!(conflicted["is_conflicted"], true); 2608 assert_eq!(conflicted["conflicts"][0]["filename"], "reef.txt"); 2609 assert_eq!(conflicted["conflicts"][0]["reason"], "patch doesn't apply"); 2610 assert_eq!(conflicted["message"], "patch cannot be applied cleanly"); 2611 2612 let broken = json_of( 2613 crate::merge::merge_check(world.state(), input("hello world\n")) 2614 .await 2615 .unwrap(), 2616 ) 2617 .await; 2618 assert_eq!(broken["is_conflicted"], true); 2619 assert!(broken["error"].as_str().is_some()); 2620 2621 assert_eq!( 2622 file_count(&repo.objects_dir()), 2623 objects_before, 2624 "merge check must write nothing into the object database" 2625 ); 2626 assert_eq!(main_tip(&world, &repo_did), base); 2627 } 2628} 2629 2630mod fork_endpoints { 2631 use super::*; 2632 2633 use knot_git::{EntryKind, Identity, NewCommit, RefUpdate, Repo, StagedAction, StagedChange}; 2634 use knot_runtime::HttpTransport; 2635 use knot_types::{Oid, RefName, UnixSeconds}; 2636 2637 const EMPTY_TREE: &str = "4b825dc642cb6eb9a060e54bf8d69288fbee4904"; 2638 2639 fn ident(time: i64) -> Identity { 2640 Identity { 2641 name: AuthorName::new("nel"), 2642 email: Email::new("nel@oyster.cafe"), 2643 time: UnixSeconds::new(time), 2644 offset_seconds: 0, 2645 } 2646 } 2647 2648 fn put_commit(repo: &Repo, parent: Option<Oid>, path: &str, content: &str, time: i64) -> Oid { 2649 let base_tree = match parent { 2650 Some(parent) => repo.find_commit(parent).unwrap().tree, 2651 None => Oid::from_hex(EMPTY_TREE).unwrap(), 2652 }; 2653 let staged = vec![StagedChange { 2654 path: knot_types::RepoPath::new(path).unwrap(), 2655 action: StagedAction::Put { 2656 content: content.as_bytes().to_vec(), 2657 kind: EntryKind::Blob, 2658 }, 2659 }]; 2660 let tree = repo.write_staged_tree(base_tree, &staged).unwrap(); 2661 repo.write_commit(&NewCommit { 2662 tree, 2663 parents: parent.into_iter().collect(), 2664 author: ident(time), 2665 committer: ident(time), 2666 message: format!("put {path}"), 2667 extra_headers: Vec::new(), 2668 }) 2669 .unwrap() 2670 } 2671 2672 fn advance(repo: &Repo, branch: &RefName, path: &str, content: &str, time: i64) -> Oid { 2673 let old = repo.find_ref(branch).unwrap(); 2674 let new = put_commit(repo, old, path, content, time); 2675 let update = match old { 2676 Some(old) => RefUpdate::Update { 2677 name: branch.clone(), 2678 old, 2679 new, 2680 }, 2681 None => RefUpdate::Create { 2682 name: branch.clone(), 2683 new, 2684 }, 2685 }; 2686 repo.update_ref(&update).unwrap(); 2687 new 2688 } 2689 2690 fn main_ref() -> RefName { 2691 RefName::new("refs/heads/main").unwrap() 2692 } 2693 2694 fn member_did() -> OwnerDid { 2695 OwnerDid::new(format!("did:web:{MEMBER_HOST}")).unwrap() 2696 } 2697 2698 fn source_url(rkey: &str) -> String { 2699 format!("https://{KNOT_HOST}/did:web:{MEMBER_HOST}/{rkey}") 2700 } 2701 2702 async fn fork_repo(world: &World, source: &str, rkey: &str) -> RepoDid { 2703 assert_eq!( 2704 as_member( 2705 world, 2706 crate::repos::create_repo, 2707 CREATE, 2708 json!({ "rkey": rkey, "name": rkey, "source": source }) 2709 ) 2710 .await 2711 .status(), 2712 StatusCode::OK 2713 ); 2714 match world 2715 .state 2716 .index 2717 .resolve_repo(&member_did(), &RepoRkey::new(rkey).unwrap()) 2718 { 2719 Resolved::Ready(Some(did)) => did, 2720 other => panic!("fork {rkey} wasn't registered: {other:?}"), 2721 } 2722 } 2723 2724 struct ForkWorld { 2725 world: World, 2726 source_did: RepoDid, 2727 fork_did: RepoDid, 2728 tip: Oid, 2729 } 2730 2731 async fn forked_world() -> ForkWorld { 2732 let world = World::new(); 2733 add_member_helper(&world).await; 2734 let source_did = create_repo_helper(&world, "kelp").await; 2735 let source = world.layout.open(&source_did).unwrap(); 2736 advance(&source, &main_ref(), "reef.txt", "kelp forest\n", 1_000); 2737 let tip = advance(&source, &main_ref(), "tide.txt", "rock pool\n", 1_001); 2738 [ 2739 "refs/tags/v1", 2740 "refs/cobs/sh.tangled.repo.collaborator/limpet", 2741 "refs/hidden/feature/main", 2742 ] 2743 .into_iter() 2744 .for_each(|name| { 2745 source 2746 .update_ref(&RefUpdate::Create { 2747 name: RefName::new(name).unwrap(), 2748 new: tip, 2749 }) 2750 .unwrap(); 2751 }); 2752 let fork_did = fork_repo(&world, &source_url("kelp"), "uni").await; 2753 ForkWorld { 2754 world, 2755 source_did, 2756 fork_did, 2757 tip, 2758 } 2759 } 2760 2761 async fn sync_fork(world: &World, signer: &K256Signer, host: &str, branch: &str) -> StatusCode { 2762 call( 2763 world, 2764 crate::forks::fork_sync, 2765 signer, 2766 host, 2767 "sh.tangled.repo.forkSync", 2768 json!({ 2769 "did": format!("did:web:{MEMBER_HOST}"), 2770 "name": "uni", 2771 "source": format!("at://did:web:{MEMBER_HOST}/sh.tangled.repo/kelp"), 2772 "branch": branch, 2773 }), 2774 ) 2775 .await 2776 .status() 2777 } 2778 2779 async fn track_hidden(world: &World, fork_ref: &str, remote_ref: &str) -> StatusCode { 2780 as_member( 2781 world, 2782 crate::forks::hidden_ref, 2783 "sh.tangled.repo.hiddenRef", 2784 json!({ 2785 "repo": format!("at://did:web:{MEMBER_HOST}/sh.tangled.repo/uni"), 2786 "forkRef": fork_ref, 2787 "remoteRef": remote_ref, 2788 }), 2789 ) 2790 .await 2791 .status() 2792 } 2793 2794 async fn fork_status( 2795 world: &World, 2796 branch: &str, 2797 hidden_ref: &str, 2798 ) -> (StatusCode, Option<u64>) { 2799 let response = as_member( 2800 world, 2801 crate::forks::fork_status, 2802 "sh.tangled.repo.forkStatus", 2803 json!({ 2804 "did": format!("did:web:{MEMBER_HOST}"), 2805 "name": "uni", 2806 "source": source_url("kelp"), 2807 "branch": branch, 2808 "hiddenRef": hidden_ref, 2809 }), 2810 ) 2811 .await; 2812 let status = response.status(); 2813 let value = json_of(response).await; 2814 (status, value["status"].as_u64()) 2815 } 2816 2817 #[tokio::test] 2818 async fn a_member_forks_a_repo_hosted_on_this_knot() { 2819 let setup = forked_world().await; 2820 assert_ne!(setup.fork_did, setup.source_did); 2821 2822 let fork = setup.world.layout.open(&setup.fork_did).unwrap(); 2823 assert_eq!(fork.find_ref(&main_ref()).unwrap(), Some(setup.tip)); 2824 assert_eq!( 2825 fork.find_ref(&RefName::new("refs/tags/v1").unwrap()) 2826 .unwrap(), 2827 Some(setup.tip) 2828 ); 2829 assert_eq!(fork.default_branch().unwrap().as_str(), "refs/heads/main"); 2830 assert_eq!( 2831 fork.origin_url().as_deref(), 2832 Some(source_url("kelp").as_str()) 2833 ); 2834 assert!( 2835 fork.references().unwrap().iter().all(|record| { 2836 !record.name.as_str().starts_with("refs/cobs/") 2837 && !record.name.as_str().starts_with("refs/hidden/") 2838 }), 2839 "a fork must copy only heads and tags, never cob or hidden refs" 2840 ); 2841 2842 let entry = fork 2843 .entry_at(setup.tip, &knot_types::RepoPath::new("tide.txt").unwrap()) 2844 .unwrap() 2845 .unwrap(); 2846 assert_eq!(fork.read_blob(entry.oid).unwrap(), b"rock pool\n"); 2847 } 2848 2849 #[tokio::test] 2850 async fn forking_a_source_this_knot_does_not_host_is_not_found() { 2851 let world = World::new(); 2852 add_member_helper(&world).await; 2853 assert_eq!( 2854 as_member(&world, crate::repos::create_repo, CREATE, json!({ "rkey": "uni", "name": "uni", "source": format!("https://{KNOT_HOST}/did:plc:whelk/ghost") })).await.status(), 2855 StatusCode::NOT_FOUND 2856 ); 2857 assert!(matches!( 2858 world 2859 .state 2860 .index 2861 .resolve_repo(&member_did(), &RepoRkey::new("uni").unwrap()), 2862 Resolved::Ready(None) 2863 )); 2864 } 2865 2866 #[tokio::test] 2867 async fn fork_sync_lifecycle() { 2868 let setup = forked_world().await; 2869 let source = setup.world.layout.open(&setup.source_did).unwrap(); 2870 let new_tip = advance(&source, &main_ref(), "spray.txt", "salt\n", 1_002); 2871 2872 assert_eq!( 2873 sync_fork(&setup.world, &setup.world.member, MEMBER_HOST, "main").await, 2874 StatusCode::OK 2875 ); 2876 let fork = setup.world.layout.open(&setup.fork_did).unwrap(); 2877 assert_eq!(fork.find_ref(&main_ref()).unwrap(), Some(new_tip)); 2878 2879 let event = only_git_event(&setup.world); 2880 assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); 2881 assert_eq!(event.payload["repo"], setup.fork_did.to_string()); 2882 assert_eq!(event.payload["ref"], "refs/heads/main"); 2883 assert_eq!(event.payload["oldSha"], setup.tip.to_string()); 2884 assert_eq!(event.payload["newSha"], new_tip.to_string()); 2885 assert_eq!( 2886 event.payload["committerDid"], 2887 account(MEMBER_HOST).to_string() 2888 ); 2889 2890 assert_eq!( 2891 sync_fork(&setup.world, &setup.world.member, MEMBER_HOST, "main").await, 2892 StatusCode::OK, 2893 "an up-to-date sync is a no-op" 2894 ); 2895 assert_eq!( 2896 git_events(&setup.world).len(), 2897 1, 2898 "an up-to-date sync emits no further event" 2899 ); 2900 2901 assert_eq!( 2902 sync_fork(&setup.world, &setup.world.stranger, STRANGER_HOST, "main").await, 2903 StatusCode::FORBIDDEN, 2904 "a stranger cannot sync a fork" 2905 ); 2906 assert_eq!( 2907 sync_fork(&setup.world, &setup.world.member, MEMBER_HOST, "driftwood").await, 2908 StatusCode::NOT_FOUND, 2909 "syncing a branch the upstream lacks isn't found" 2910 ); 2911 } 2912 2913 #[tokio::test] 2914 async fn hidden_ref_tracks_the_upstream_branch_and_stays_hidden() { 2915 let setup = forked_world().await; 2916 let source = setup.world.layout.open(&setup.source_did).unwrap(); 2917 let new_tip = advance(&source, &main_ref(), "spray.txt", "salt\n", 1_002); 2918 2919 assert_eq!( 2920 track_hidden(&setup.world, "feature", "main").await, 2921 StatusCode::OK 2922 ); 2923 2924 let fork = setup.world.layout.open(&setup.fork_did).unwrap(); 2925 let hidden = RefName::new("refs/hidden/feature/main").unwrap(); 2926 assert_eq!(fork.find_ref(&hidden).unwrap(), Some(new_tip)); 2927 assert!( 2928 fork.advertised_refs() 2929 .unwrap() 2930 .iter() 2931 .all(|record| !record.name.as_str().starts_with("refs/hidden/")), 2932 "a hidden ref must stay out of the public advertisement" 2933 ); 2934 assert_eq!( 2935 track_hidden(&setup.world, "feature", "main").await, 2936 StatusCode::OK, 2937 "tracking an already-tracked ref is idempotent" 2938 ); 2939 } 2940 2941 #[tokio::test] 2942 async fn fork_status_reports_up_to_date_fast_forwardable_and_conflict() { 2943 let setup = forked_world().await; 2944 assert_eq!( 2945 track_hidden(&setup.world, "feature", "main").await, 2946 StatusCode::OK 2947 ); 2948 assert_eq!( 2949 fork_status(&setup.world, "main", "refs/hidden/feature/main").await, 2950 (StatusCode::OK, Some(0)) 2951 ); 2952 2953 let source = setup.world.layout.open(&setup.source_did).unwrap(); 2954 advance(&source, &main_ref(), "spray.txt", "salt\n", 1_002); 2955 assert_eq!( 2956 track_hidden(&setup.world, "feature", "main").await, 2957 StatusCode::OK 2958 ); 2959 assert_eq!( 2960 fork_status(&setup.world, "main", "refs/hidden/feature/main").await, 2961 (StatusCode::OK, Some(1)) 2962 ); 2963 2964 let fork = setup.world.layout.open(&setup.fork_did).unwrap(); 2965 advance(&fork, &main_ref(), "wreck.txt", "barnacle\n", 1_003); 2966 assert_eq!( 2967 fork_status(&setup.world, "main", "refs/hidden/feature/main").await, 2968 (StatusCode::OK, Some(2)) 2969 ); 2970 2971 assert_eq!( 2972 fork_status(&setup.world, "main", "refs/hidden/ghost/main").await, 2973 (StatusCode::BAD_REQUEST, None), 2974 "an unresolvable revision is an invalid request" 2975 ); 2976 } 2977 2978 #[tokio::test] 2979 async fn fork_status_reports_up_to_date_when_the_fork_is_ahead() { 2980 let setup = forked_world().await; 2981 assert_eq!( 2982 track_hidden(&setup.world, "feature", "main").await, 2983 StatusCode::OK 2984 ); 2985 let fork = setup.world.layout.open(&setup.fork_did).unwrap(); 2986 advance(&fork, &main_ref(), "wreck.txt", "barnacle\n", 1_003); 2987 assert_eq!( 2988 fork_status(&setup.world, "main", "refs/hidden/feature/main").await, 2989 (StatusCode::OK, Some(0)) 2990 ); 2991 } 2992 2993 #[tokio::test] 2994 async fn a_fork_from_a_remote_knot_fetches_over_http() { 2995 let upstream_dir = tempfile::tempdir().unwrap(); 2996 let upstream_path = upstream_dir.path().join("uni.git"); 2997 let upstream = Repo::create(&upstream_path).unwrap(); 2998 upstream 2999 .set_head(&RefName::new("refs/heads/main").unwrap()) 3000 .unwrap(); 3001 advance(&upstream, &main_ref(), "reef.txt", "kelp forest\n", 1_000); 3002 let tip = advance(&upstream, &main_ref(), "tide.txt", "rock pool\n", 1_001); 3003 3004 let served = upstream_path.clone(); 3005 let git_http: Arc<dyn HttpTransport> = 3006 Arc::new(FakeHttp::new(move |request: &HttpRequest| { 3007 let _keep = &upstream_dir; 3008 let repo = Repo::open(&served).unwrap(); 3009 let body = if request.url.path().ends_with("/info/refs") { 3010 knot_pack::advertise_upload(&repo).unwrap() 3011 } else { 3012 knot_pack::upload_pack(&repo, request.body.as_deref().unwrap_or_default()) 3013 .unwrap() 3014 }; 3015 Ok(HttpResponse { 3016 status: StatusCode::OK, 3017 headers: http::HeaderMap::new(), 3018 body: body.into(), 3019 }) 3020 })); 3021 3022 let world = World::with_git_http(git_http); 3023 add_member_helper(&world).await; 3024 let remote = "https://barnacle.nel.pet/did:plc:squid/uni"; 3025 let fork_did = fork_repo(&world, remote, "uni").await; 3026 let fork = world.layout.open(&fork_did).unwrap(); 3027 assert_eq!(fork.find_ref(&main_ref()).unwrap(), Some(tip)); 3028 assert_eq!(fork.origin_url().as_deref(), Some(remote)); 3029 3030 let new_tip = advance( 3031 &Repo::open(&upstream_path).unwrap(), 3032 &main_ref(), 3033 "spray.txt", 3034 "salt\n", 3035 1_002, 3036 ); 3037 assert_eq!( 3038 sync_fork(&world, &world.member, MEMBER_HOST, "main").await, 3039 StatusCode::OK 3040 ); 3041 let fork = world.layout.open(&fork_did).unwrap(); 3042 assert_eq!(fork.find_ref(&main_ref()).unwrap(), Some(new_tip)); 3043 } 3044} 3045 3046mod legacy_admin_route { 3047 use super::*; 3048 use crate::legacy_admin::{ADD_MEMBER_ROUTE, LegacyAdminSecret}; 3049 use tower::ServiceExt; 3050 3051 const SECRET: &str = "nekomilk2"; 3052 3053 async fn call(router: &axum::Router, user: &str, password: &str, subject: &str) -> StatusCode { 3054 let encoded = base64::engine::general_purpose::STANDARD 3055 .encode(format!("{user}:{password}").as_bytes()); 3056 let request = http::Request::builder() 3057 .method("POST") 3058 .uri(ADD_MEMBER_ROUTE) 3059 .header(AUTHORIZATION, format!("Basic {encoded}")) 3060 .header(http::header::CONTENT_TYPE, "application/json") 3061 .body(axum::body::Body::from( 3062 json!({ "subject": subject }).to_string(), 3063 )) 3064 .unwrap(); 3065 router.clone().oneshot(request).await.unwrap().status() 3066 } 3067 3068 #[tokio::test] 3069 async fn the_legacy_route_admits_a_member_only_with_the_configured_credentials() { 3070 let world = World::new(); 3071 let router = crate::router(Arc::clone(&world.state)).merge(crate::legacy_admin::router( 3072 Arc::clone(&world.state), 3073 LegacyAdminSecret::new(SECRET).unwrap(), 3074 )); 3075 let subject = format!("did:web:{MEMBER_HOST}"); 3076 3077 let version = http::Request::builder() 3078 .method("GET") 3079 .uri(crate::service::VERSION_ROUTE) 3080 .body(axum::body::Body::empty()) 3081 .unwrap(); 3082 assert_eq!( 3083 router.clone().oneshot(version).await.unwrap().status(), 3084 StatusCode::OK, 3085 "merging the legacy route leaves the xrpc routes reachable" 3086 ); 3087 3088 let refused = futures::future::join_all( 3089 [("admin", "nope"), ("root", SECRET), ("admin", "")] 3090 .map(|(user, password)| call(&router, user, password, &subject)), 3091 ) 3092 .await; 3093 assert!( 3094 refused 3095 .iter() 3096 .all(|status| *status == StatusCode::UNAUTHORIZED), 3097 "the knot refuses a wrong user or secret, got {refused:?}" 3098 ); 3099 assert_eq!( 3100 call( 3101 &router, 3102 "admin", 3103 SECRET, 3104 &"n".repeat(world.state.byte_limits.body.get() + 1) 3105 ) 3106 .await, 3107 StatusCode::PAYLOAD_TOO_LARGE 3108 ); 3109 assert_eq!( 3110 world.state.index.is_member(&account(MEMBER_HOST)), 3111 Resolved::Ready(false), 3112 "a refused call grants nothing" 3113 ); 3114 3115 assert_eq!( 3116 call(&router, "admin", SECRET, &subject).await, 3117 StatusCode::OK 3118 ); 3119 assert_eq!( 3120 world.state.index.is_member(&account(MEMBER_HOST)), 3121 Resolved::Ready(true) 3122 ); 3123 let added = last_event(&world, "sh.tangled.knot.memberUpdate"); 3124 assert_eq!(added.payload["op"], "add"); 3125 assert_eq!(added.payload["subject"], account(MEMBER_HOST).to_string()); 3126 let Resolved::Ready(members) = world.state.index.member_entries() else { 3127 panic!("the member roster is warm in this test"); 3128 }; 3129 assert_eq!( 3130 members 3131 .iter() 3132 .find(|grant| grant.subject == account(MEMBER_HOST)) 3133 .expect("the member is in the roster") 3134 .added_by, 3135 world.state.service_owner, 3136 "the legacy grant records the service owner as the granter" 3137 ); 3138 3139 let baseline = event_count(&world); 3140 assert_eq!( 3141 call(&router, "admin", SECRET, &subject).await, 3142 StatusCode::OK, 3143 "the legacy route is idempotent, matching the Go knot" 3144 ); 3145 assert_eq!( 3146 event_count(&world), 3147 baseline, 3148 "re-adding an existing member emits no event" 3149 ); 3150 } 3151 3152 #[tokio::test] 3153 async fn the_legacy_route_sheds_a_pre_auth_flood_from_one_peer() { 3154 use axum::extract::ConnectInfo; 3155 use std::net::SocketAddr; 3156 3157 let world = World::new(); 3158 let router = crate::legacy_admin::router( 3159 Arc::clone(&world.state), 3160 LegacyAdminSecret::new(SECRET).unwrap(), 3161 ); 3162 let peer = SocketAddr::from(([203, 0, 113, 9], 5555)); 3163 3164 let statuses: Vec<StatusCode> = futures::stream::iter(0..22) 3165 .then(|_| { 3166 let router = router.clone(); 3167 async move { 3168 let mut request = http::Request::builder() 3169 .method("POST") 3170 .uri(ADD_MEMBER_ROUTE) 3171 .body(axum::body::Body::empty()) 3172 .unwrap(); 3173 request.extensions_mut().insert(ConnectInfo(peer)); 3174 router.oneshot(request).await.unwrap().status() 3175 } 3176 }) 3177 .collect() 3178 .await; 3179 3180 assert!( 3181 statuses[..20] 3182 .iter() 3183 .all(|status| *status == StatusCode::UNAUTHORIZED), 3184 "the knot admits the per-peer burst and then fails it on the missing credentials, got {statuses:?}" 3185 ); 3186 assert!( 3187 statuses[20..] 3188 .iter() 3189 .all(|status| *status == StatusCode::TOO_MANY_REQUESTS), 3190 "past the burst the knot sheds the guess flood before it reaches the secret comparison, got {statuses:?}" 3191 ); 3192 } 3193}