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