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