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