Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188use std::collections::{HashMap, HashSet};use std::path::{Path, PathBuf};use std::process::Command;use std::sync::atomic::{AtomicUsize, Ordering};use std::sync::{Arc, Mutex};
use futures::stream::StreamExt;use knot_atproto::Atproto;use knot_cob::{CobHome, CobStore};use knot_cobs::{CollaboratorsChange, Grant, MembersChange, Registration, RegistryChange};use knot_fixtures::UNSHARED_SSH;use knot_git::{ArchiveLimit, Layout, Repo};use knot_index::{Index, Resolved};use knot_pack::MaxWireBytes;use knot_runtime::{ FakeDns, FakeHttp, HttpResponse, K256Signer, ManualClock, SeededEntropy, Signer, UnixMicros,};use knot_types::{ AccountDid, KnotId, Oid, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, UnixSeconds,};use tempfile::TempDir;use tokio::net::TcpListener;use url::Url;
const REPO_DID: &str = "did:plc:squid";const REPO_NAME: &str = "anemone";const OWNER_DID: &str = "did:plc:nel";const TID_REPO_DID: &str = "did:plc:limpet";const TID_RKEY: &str = "3mizfnpxii522";const TID_REPO_NAME: &str = "periwinkle.cloud";const PDS_HOST: &str = "pds.oyster.cafe";
fn git(cwd: &Path, env: &[(&str, &str)], args: &[&str]) -> (bool, String) { let mut command = knot_fixtures::command(cwd); command.args(args); env.iter().for_each(|(key, value)| { command.env(key, value); }); let out = command.output().expect("git runs"); let combined = format!( "{}{}", String::from_utf8_lossy(&out.stdout), String::from_utf8_lossy(&out.stderr) ); (out.status.success(), combined)}
fn keygen(dir: &Path, name: &str) -> (String, String) { let path = dir.join(name); let out = Command::new("ssh-keygen") .args([ "-t", "ed25519", "-N", "", "-C", "nel@oyster.cafe", "-f", path.to_str().unwrap(), ]) .output() .expect("ssh-keygen runs"); assert!( out.status.success(), "ssh-keygen failed: {}", String::from_utf8_lossy(&out.stderr) ); let public_line = std::fs::read_to_string(dir.join(format!("{name}.pub"))) .unwrap() .trim() .to_string(); (path.to_str().unwrap().to_string(), public_line)}
fn did_document(signer: &K256Signer, did: &str, pds: &str) -> Vec<u8> { let multikey = knot_types::crypto::multikey(0xe7, signer.public_key().as_bytes()); serde_json::to_vec(&serde_json::json!({ "id": did, "alsoKnownAs": ["at://nel.pet"], "verificationMethod": [{ "id": format!("{did}#atproto"), "type": "Multikey", "controller": did, "publicKeyMultibase": multikey }], "service": [{ "id": "#atproto_pds", "type": "AtprotoPersonalDataServer", "serviceEndpoint": pds }] })) .unwrap()}
fn list_records_body(lines: &[&str]) -> Vec<u8> { let records: Vec<_> = lines .iter() .map(|line| { serde_json::json!({ "value": { "$type": "sh.tangled.publicKey", "key": line, "name": "laptop", "createdAt": "2026-06-08T00:00:00Z" } }) }) .collect(); serde_json::to_vec(&serde_json::json!({ "records": records })).unwrap()}
fn ok_body(body: Vec<u8>) -> HttpResponse { HttpResponse { status: http::StatusCode::OK, headers: http::HeaderMap::new(), body: bytes::Bytes::from(body), }}
fn fake_dns() -> impl knot_runtime::DnsTxtResolver { FakeDns::new(|name: &str| { Ok(match name { "_atproto.nel.pet" => vec![format!("did={OWNER_DID}")], _ => Vec::new(), }) })}
fn not_found() -> HttpResponse { HttpResponse { status: http::StatusCode::NOT_FOUND, headers: http::HeaderMap::new(), body: bytes::Bytes::new(), }}
fn forever() -> knot_index::KeyLease { knot_index::KeyTtl::from_secs(u32::MAX.into()).lease_from(UnixSeconds::new(0))}
fn server_error() -> HttpResponse { HttpResponse { status: http::StatusCode::INTERNAL_SERVER_ERROR, headers: http::HeaderMap::new(), body: bytes::Bytes::new(), }}
fn fake_http(published_line: String) -> impl knot_runtime::HttpTransport { let signer = K256Signer::generate(&SeededEntropy::new(1)); let pds = format!("https://{PDS_HOST}"); FakeHttp::new(move |request| { let host = request.url.host_str().unwrap_or_default().to_string(); let path = request.url.path().to_string(); let body = if host == PDS_HOST { list_records_body(&[&published_line]) } else if path.ends_with(REPO_DID) { did_document(&signer, REPO_DID, &pds) } else if path.ends_with(TID_REPO_DID) { did_document(&signer, TID_REPO_DID, &pds) } else if path.ends_with(OWNER_DID) { did_document(&signer, OWNER_DID, &pds) } else { return Ok(not_found()); }; Ok(ok_body(body)) })}
#[derive(Default, Clone)]struct Accounts { identities: Arc<Mutex<HashMap<String, Vec<String>>>>, unreachable: Arc<Mutex<HashSet<String>>>, listings: Arc<AtomicUsize>,}
impl Accounts { fn publishing(identities: HashMap<String, Vec<String>>) -> Self { Self { identities: Arc::new(Mutex::new(identities)), ..Self::default() } }
fn unreachable(self, dids: HashSet<String>) -> Self { *self.unreachable.lock().unwrap() = dids; self }
fn restore(&self, did: &str) { self.unreachable.lock().unwrap().remove(did); }
fn publish(&self, did: &str, line: String) { self.identities .lock() .unwrap() .entry(did.to_string()) .or_default() .push(line); }
fn published_by(&self, did: &str) -> Vec<String> { self.identities .lock() .unwrap() .get(did) .cloned() .unwrap_or_default() }
fn listings(&self) -> usize { self.listings.load(Ordering::SeqCst) }}
fn multi_http(accounts: Accounts) -> impl knot_runtime::HttpTransport { let signer = K256Signer::generate(&SeededEntropy::new(77)); FakeHttp::new(move |request| { let host = request.url.host_str().unwrap_or_default().to_string(); if host == "plc.directory" { let did = request.url.path().trim_start_matches('/').to_string(); return Ok(ok_body(did_document(&signer, &did, "https://pds.test"))); } if host == "pds.test" { let repo = request .url .query_pairs() .find(|(key, _)| key == "repo") .map(|(_, value)| value.into_owned()) .unwrap_or_default(); accounts.listings.fetch_add(1, Ordering::SeqCst); if accounts.unreachable.lock().unwrap().contains(&repo) { return Ok(server_error()); } let lines = accounts.published_by(&repo); let refs: Vec<&str> = lines.iter().map(String::as_str).collect(); return Ok(ok_body(list_records_body(&refs))); } Ok(not_found()) })}
fn actor_for_seed(seed: u64) -> knot_types::ActorId { knot_types::ActorId::from_secp256k1( K256Signer::generate(&SeededEntropy::new(seed)) .public_key() .as_bytes(), )}
struct Server { _scan: TempDir, layout: Layout, repo_did: RepoDid, port: u16,}
async fn spawn_server( published_line: String, max_pack_bytes: MaxWireBytes,) -> (Server, Arc<Index>) { spawn_server_with(published_line, max_pack_bytes, true).await}
async fn spawn_server_with( published_line: String, max_pack_bytes: MaxWireBytes, warm: bool,) -> (Server, Arc<Index>) { let (server, index, _, _, _) = spawn_server_core( published_line, max_pack_bytes, ArchiveLimit::default(), warm, None, false, ) .await; (server, index)}
async fn spawn_server_core( published_line: String, max_pack_bytes: MaxWireBytes, archive_limit: ArchiveLimit, warm: bool, lfs: Option<knot_lfs::LfsHandle>, projection: bool,) -> ( Server, Arc<Index>, tokio_util::sync::CancellationToken, tokio::task::JoinHandle<()>, Option<Arc<knot_events::EventLog<ManualClock>>>,) { let scan = tempfile::tempdir().unwrap(); let meta_path = scan.path().join("meta"); Repo::create(&meta_path).unwrap(); let layout = Layout::new(scan.path().join("repos")); let repo_did = RepoDid::new(REPO_DID).unwrap(); layout.create(&repo_did).unwrap();
let signer = K256Signer::generate(&SeededEntropy::new(2)); let meta = Repo::open(&meta_path).unwrap(); let store = CobStore::new(&meta); let home = CobHome::from(&KnotId::new("did:web:nel.pet").unwrap()); let registry = store .create( &home, &RegistryChange::Register(Registration { owner: OwnerDid::new(OWNER_DID).unwrap(), rkey: RepoRkey::new(REPO_NAME).unwrap(), name: RepoName::new(REPO_NAME).unwrap(), repo: repo_did.clone(), created_at: UnixSeconds::new(1), policy: knot_types::ContributionPolicy::default(), }), &signer, UnixSeconds::new(1), ) .unwrap() .object;
let tid_repo_did = RepoDid::new(TID_REPO_DID).unwrap(); layout.create(&tid_repo_did).unwrap(); store .update( &home, registry, &RegistryChange::Register(Registration { owner: OwnerDid::new(OWNER_DID).unwrap(), rkey: RepoRkey::new(TID_RKEY).unwrap(), name: RepoName::new(TID_REPO_NAME).unwrap(), repo: tid_repo_did, created_at: UnixSeconds::new(2), policy: knot_types::ContributionPolicy::default(), }), &signer, UnixSeconds::new(2), ) .unwrap();
let index = Arc::new(Index::new(meta_path, layout.clone())); if warm { index.rebuild().unwrap(); }
let atproto = Arc::new( Atproto::new( fake_http(published_line), ManualClock::new(UnixMicros::new(1_000_000_000)), KnotId::new("did:web:nel.pet").unwrap(), knot_atproto::PlcDirectory::new(Url::parse("https://plc.directory/").unwrap()).unwrap(), ) .with_dns(Arc::new(fake_dns())), );
let key_dir = scan.path().join("hostkey"); std::fs::create_dir_all(&key_dir).unwrap(); let host_key = knot_ssh::load_or_create_host_key(&key_dir.join("host")).unwrap();
let (xrpc_state, firehose) = if projection { let knot = KnotId::new("did:web:nel.pet").unwrap(); let secrets = Arc::new( knot_secrets::SealedStore::open( scan.path().join("keys.sealed"), &knot_secrets::MasterKey::new([7u8; 32]).unwrap(), Box::new(knot_runtime::OsEntropy), ) .unwrap(), ); secrets.ensure(&knot).unwrap(); let state = Arc::new(knot_xrpc::XrpcState { imports: knot_xrpc::imports::ImportQueue::channel().0, layout: layout.clone(), index: Arc::clone(&index), atproto: Arc::clone(&atproto), secrets, entropy: Arc::new(knot_runtime::OsEntropy), invite_codes: Default::default(), ci_logs: None, admins: std::collections::BTreeSet::new(), admission: knot_types::AdmissionPolicy::Closed, contribution_policy: knot_types::ContributionPolicy::Anyone, writes: Arc::new(knot_resource::SubjectLimiter::new( knot_resource::RateLimit { burst: knot_xrpc::Burst::literal(1_000), refill: knot_xrpc::RefillMicros::literal(1_000), }, )), social_budget: knot_resource::SocialBudgetBytes::new(1 << 28), knot_did: knot, knot_hostname: knot_types::KnotHostname::new("knot.test").unwrap(), meta_path: scan.path().join("meta"), knot_service_url: knot_types::KnotServiceUrl::new("https://knot.test".to_string()) .unwrap(), limiter: Arc::new(knot_xrpc::PreAuthLimiter::with_config( knot_xrpc::LimitConfig::default(), )), cob_locks: Arc::new(knot_xrpc::CobLocks::default()), reservations: Arc::new(knot_xrpc::Reservations::new( knot_xrpc::ReservationTtl::new(1_000_000), knot_xrpc::PerActorQuota::new(16), knot_xrpc::GlobalQuota::new(16), )), proxy_trust: knot_types::ProxyTrust::default(), committer: knot_xrpc::Committer { name: knot_types::AuthorName::new("Tangled"), email: knot_types::Email::new("noreply@tangled.sh"), }, byte_limits: knot_xrpc::ByteLimits::default(), budgets: knot_xrpc::Budgets::default(), git_http: Arc::new(knot_runtime::FakeHttp::new( |_request: &knot_runtime::HttpRequest| Ok(not_found()), )), pack_limits: knot_pack::PackLimits::default(), service_owner: AccountDid::new(OWNER_DID).unwrap(), subscriber_gate: Arc::new(knot_events::SubscriberGate::new( knot_events::GlobalSubscriberLimit::new(16), knot_events::PerPeerSubscriberLimit::new(4), )), maintenance: knot_maintenance::MaintenanceHandle::disabled(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), slots: knot_resource::Slots::testing(8), lfs: None, blobs: Default::default(), catalog: Arc::new(knot_messages::Catalog::defaults()), firehose: Arc::new(knot_events::EventLog::new( ManualClock::new(UnixMicros::new(1_000_000_000)), knot_events::ReplayBounds::new( knot_events::ReplayEvents::new(4096).unwrap(), knot_events::ReplayBytes::new(16 << 20).unwrap(), ), )), rkeys: Default::default(), }); let firehose = Arc::clone(&state.firehose); (Some(state), Some(firehose)) } else { (None, None) }; let base = knot_ssh::SshState::new(knot_ssh::SshConfig { layout: layout.clone(), index: Arc::clone(&index), atproto, knot_actor: actor_for_seed(1), hostname: knot_types::KnotHostname::new("knot.test").unwrap(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), admins: std::collections::BTreeSet::new(), admission: knot_types::AdmissionPolicy::Closed, max_pack_bytes, archive_limit, ci_logs: None, }); let base = match xrpc_state { Some(xrpc) => base.with_ref_surface(Arc::new(knot_xrpc::RefProjection { state: xrpc })), None => base, }; let state = Arc::new(match lfs { Some(handle) => base.with_lfs(handle, 2), None => base, });
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let port = listener.local_addr().unwrap().port(); let shutdown = tokio_util::sync::CancellationToken::new(); let serve_task = tokio::spawn({ let shutdown = shutdown.clone(); async move { let _ = knot_ssh::serve_drained(listener, host_key, state, shutdown).await; } });
( Server { _scan: scan, layout, repo_did, port, }, index, shutdown, serve_task, firehose, )}
fn ssh_command(key_path: &str) -> String { format!("ssh -i {key_path} {UNSHARED_SSH}")}
async fn git_ssh(cwd: &Path, key: &str, args: &[&str]) -> (bool, String) { let ssh = ssh_command(key); let cwd = cwd.to_path_buf(); let owned: Vec<String> = args.iter().map(|arg| arg.to_string()).collect(); tokio::task::spawn_blocking(move || { let argv: Vec<&str> = owned.iter().map(String::as_str).collect(); git(&cwd, &[("GIT_SSH_COMMAND", &ssh)], &argv) }) .await .unwrap()}
async fn push(work: &Path, url: &str, key: &str, refspecs: &[&str]) -> (bool, String) { let args: Vec<&str> = std::iter::once("push") .chain(std::iter::once(url)) .chain(refspecs.iter().copied()) .collect(); git_ssh(work, key, &args).await}
async fn clone(url: &str, key: &str, dest: &Path) -> (bool, String) { git_ssh( Path::new("/tmp"), key, &["clone", "-q", url, dest.to_str().unwrap()], ) .await}
fn seed_work(work: &Path) -> String { std::fs::create_dir_all(work).unwrap(); git(work, &[], &["init", "-q", "-b", "main"]); std::fs::write(work.join("README.md"), "hello over ssh\n").unwrap(); git(work, &[], &["add", "-A"]); git(work, &[], &["commit", "-q", "-m", "initial"]); let (ok, head) = git(work, &[], &["rev-parse", "HEAD"]); assert!(ok); head.trim().to_string()}
fn seed_commits(work: &Path, count: usize) { std::fs::create_dir_all(work).unwrap(); git(work, &[], &["init", "-q", "-b", "main"]); (0..count).for_each(|i| { std::fs::write(work.join("log.txt"), format!("line {i}\n")).unwrap(); git(work, &[], &["add", "-A"]); git(work, &[], &["commit", "-q", "-m", &format!("c{i}")]); });}
fn seed_cob(work: &Path, signer_seed: u64, subject: &str, home: &CobHome) -> (Oid, String, String) { let repo = Repo::open(work).unwrap(); let signer = K256Signer::generate(&SeededEntropy::new(signer_seed)); let created = CobStore::new(&repo) .create( home, &MembersChange::Add(Grant { subject: AccountDid::new(subject).unwrap(), added_by: AccountDid::new(OWNER_DID).unwrap(), created_at: UnixSeconds::new(1), }), &signer, UnixSeconds::new(1), ) .unwrap(); let cob_ref = format!( "refs/cobs/sh.tangled.knot.member/{}", created.object.oid().to_hex() ); let spec = format!("{cob_ref}:{cob_ref}"); (created.tip.oid(), cob_ref, spec)}
fn main_tip(layout: &Layout, repo: &RepoDid) -> Option<Oid> { layout .open(repo) .unwrap() .find_ref(&RefName::new("refs/heads/main").unwrap()) .unwrap()}
fn ref_names(server: &Server) -> Vec<String> { server .layout .open(&server.repo_did) .unwrap() .references() .unwrap() .iter() .map(|record| record.name.as_str().to_string()) .collect()}
struct Fixture { scratch: TempDir, server: Server, index: Arc<Index>, key_path: String, url: String, work: PathBuf,}
async fn fixture() -> Fixture { fixture_with_archive_limit(ArchiveLimit::default()).await}
async fn fixture_with_archive_limit(archive_limit: ArchiveLimit) -> Fixture { let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "client"); let (server, index, _, _, _) = spawn_server_core( public_line, MaxWireBytes::new(1 << 30), archive_limit, true, None, false, ) .await; let url = format!( "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", server.port ); let work = scratch.path().join("work"); Fixture { scratch, server, index, key_path, url, work, }}
fn fetch_main_exit(clone_dir: &Path, ssh: &str, extra_git: &[&str]) -> Option<i32> { let mut args = vec!["-k", "3", "20", "git"]; args.extend_from_slice(extra_git); args.extend_from_slice(&["fetch", "origin", "main"]); Command::new("timeout") .args(&args) .current_dir(clone_dir) .env("GIT_SSH_COMMAND", ssh) .status() .expect("timeout/git runs") .code()}
async fn incremental_fetch_exit( seed_count: usize, extra_git: &'static [&'static str],) -> Option<i32> { let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "client"); let (server, _index) = spawn_server(public_line, MaxWireBytes::new(1 << 30)).await; let url = format!( "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", server.port );
let work = scratch.path().join("work"); seed_commits(&work, seed_count); let (ok, out) = push(&work, &url, &key_path, &["main"]).await; assert!(ok, "seeding push must land:\n{out}");
let clone_dir = scratch.path().join("clone"); let (ok, out) = clone(&url, &key_path, &clone_dir).await; assert!(ok, "clone over ssh must succeed:\n{out}");
git( &work, &[], &["commit", "-q", "--allow-empty", "-m", "advance"], ); let (ok, out) = push(&work, &url, &key_path, &["main"]).await; assert!(ok, "advancing server tip must succeed:\n{out}");
let ssh = ssh_command(&key_path); let exit = tokio::task::spawn_blocking(move || fetch_main_exit(&clone_dir, &ssh, extra_git)) .await .unwrap(); drop(server); exit}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn incremental_fetch_over_ssh_completes() { let cases: [(&'static [&'static str], &str); 2] = [ ( &["-c", "protocol.version=0"], "diverged v0 fetch sends more than 32 haves and blocks on an ACK/NAK. Upload loop \ answers each have-batch flush with a NAK instead of waiting for done, so it never \ hangs", ), ( &[], "git forwards GIT_PROTOCOL over ssh, so default fetch path negotiates with the v2 loop", ), ]; futures::stream::iter(cases) .for_each(|(extra, rationale)| async move { assert_eq!( incremental_fetch_exit(50, extra).await, Some(0), "{rationale}" ); }) .await;}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_display_name_addresses_a_repo_whose_record_key_is_a_tid() { let fx = fixture().await; let head = seed_work(&fx.work); let head_oid = Oid::from_hex(&head).unwrap(); let port = fx.server.port; let target = RepoDid::new(TID_REPO_DID).unwrap(); let variants = [ format!("ssh://git@127.0.0.1:{port}/{OWNER_DID}/{TID_REPO_NAME}"), format!("ssh://git@127.0.0.1:{port}/{OWNER_DID}/{TID_REPO_NAME}.git"), format!("ssh://git@127.0.0.1:{port}/nel.pet/{TID_REPO_NAME}"), format!("ssh://git@127.0.0.1:{port}/{OWNER_DID}/{TID_RKEY}"), ]; let fx = &fx; let target = ⌖ futures::stream::iter(variants) .for_each(|url| async move { let (ok, out) = push(&fx.work, &url, &fx.key_path, &["main"]).await; assert!( ok, "a PDS-minted record key leaves the display name as the only human \ path, so {url} must resolve and push:\n{out}" ); assert_eq!( main_tip(&fx.server.layout, target), Some(head_oid), "{url}: pushed commit must be the named repository's main tip" ); }) .await;
let (ok, out) = push( &fx.work, &format!("ssh://git@127.0.0.1:{port}/{OWNER_DID}/whelk"), &fx.key_path, &["main"], ) .await; assert!( !ok, "a segment matching neither a record key nor a name stays unresolvable:\n{out}" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn addressing_variants_land() { let fx = fixture().await; let head = seed_work(&fx.work); let head_oid = Oid::from_hex(&head).unwrap(); let port = fx.server.port; let variants = [ format!("ssh://git@127.0.0.1:{port}/{OWNER_DID}/{REPO_NAME}"), format!("ssh://git@127.0.0.1:{port}/{OWNER_DID}/{REPO_NAME}.git"), format!("ssh://git@127.0.0.1:{port}/nel.pet/{REPO_NAME}"), format!("ssh://git@127.0.0.1:{port}/{REPO_DID}"), ]; let fx = &fx; futures::stream::iter(variants) .for_each(|url| async move { let (ok, out) = push(&fx.work, &url, &fx.key_path, &["main"]).await; assert!(ok, "addressing {url} must resolve and push:\n{out}"); assert_eq!( main_tip(&fx.server.layout, &fx.server.repo_did), Some(head_oid), "{url}: pushed commit must be the repository's main tip" ); }) .await;}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_push_while_the_index_is_warming_is_refused() { let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "client"); let (server, _index) = spawn_server_with(public_line, MaxWireBytes::new(1 << 30), false).await; let url = format!( "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", server.port );
let work = scratch.path().join("work"); seed_work(&work); let (ok, out) = push(&work, &url, &key_path, &["main"]).await; assert!( !ok, "warming index must fail closed at the SSH boundary:\n{out}" ); assert_eq!( main_tip(&server.layout, &server.repo_did), None, "no ref lands while index is warming" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn unresolvable_targets_refused() { let fx = fixture().await; seed_work(&fx.work); let port = fx.server.port;
let bad_name = format!("ssh://git@127.0.0.1:{port}/{OWNER_DID}/conch"); let (ok, out) = push(&fx.work, &bad_name, &fx.key_path, &["main"]).await; assert!( !ok, "owner/name with no registry entry must be rejected, not silently routed:\n{out}" );
fx.server .layout .create(&RepoDid::new("did:plc:clam").unwrap()) .unwrap(); let ghost_url = format!("ssh://git@127.0.0.1:{port}/did:plc:clam"); let dest = fx.scratch.path().join("ghost"); let (ok, out) = clone(&ghost_url, &fx.key_path, &dest).await; assert!( !ok, "repo present on disk but absent from registry mustn't be served by bare DID:\n{out}" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn an_authorized_push_over_ssh_succeeds_and_a_clone_reads_it_back() { let fx = fixture().await; let head = seed_work(&fx.work);
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "authorized push over ssh must succeed:\n{out}"); assert_eq!( main_tip(&fx.server.layout, &fx.server.repo_did), Some(Oid::from_hex(&head).unwrap()), "pushed commit must be the repository's main tip" );
let clone_dir = fx.scratch.path().join("clone"); let (ok, out) = clone(&fx.url, &fx.key_path, &clone_dir).await; assert!(ok, "clone over ssh must succeed:\n{out}"); assert!( clone_dir.join("README.md").exists(), "clone must check out the pushed file" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_client_requesting_ssh_compression_clones_an_incompressible_pack() { let fx = fixture().await; std::fs::create_dir_all(&fx.work).unwrap(); git(&fx.work, &[], &["init", "-q", "-b", "main"]); let mut state = 0x9e3779b97f4a7c15u64; let payload: Vec<u8> = std::iter::repeat_with(|| { state ^= state << 13; state ^= state >> 7; state ^= state << 17; state.to_le_bytes() }) .take(32 * 1024) .flatten() .collect(); std::fs::write(fx.work.join("noise.bin"), &payload).unwrap(); git(&fx.work, &[], &["add", "-A"]); git(&fx.work, &[], &["commit", "-q", "-m", "noise"]);
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "push must succeed:\n{out}");
let dest = fx.scratch.path().join("compressed-clone"); let ssh = format!("{} -o Compression=yes", ssh_command(&fx.key_path)); let url = fx.url.clone(); let dest_arg = dest.to_str().unwrap().to_string(); let (ok, out) = tokio::task::spawn_blocking(move || { git( Path::new("/tmp"), &[("GIT_SSH_COMMAND", &ssh)], &["clone", "-q", &url, &dest_arg], ) }) .await .unwrap(); assert!( ok, "clone with ssh compression requested must succeed:\n{out}" ); assert_eq!(std::fs::read(dest.join("noise.bin")).unwrap(), payload);}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn an_oversized_push_is_refused_at_the_ssh_boundary() { let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "client"); let (server, _index) = spawn_server(public_line, MaxWireBytes::new(64)).await; let url = format!( "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", server.port );
let work = scratch.path().join("work"); seed_work(&work); let (ok, out) = push(&work, &url, &key_path, &["main"]).await; assert!( !ok, "push larger than the configured limit must be refused:\n{out}" ); assert!( ref_names(&server).is_empty(), "oversized push mustn't land any ref" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn an_up_to_date_push_over_ssh_is_accepted() { let fx = fixture().await; seed_work(&fx.work);
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "first push must land:\n{out}");
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!( ok, "up-to-date no-op push must succeed instead of failing with a stream error:\n{out}" ); assert!( out.contains("up-to-date") || out.contains("up to date"), "git must report branch is up to date:\n{out}" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_denied_push_over_ssh_leaves_no_objects_in_the_live_odb() { let scratch = tempfile::tempdir().unwrap(); let (_registered_path, registered_line) = keygen(scratch.path(), "registered"); let (attacker_path, _attacker_line) = keygen(scratch.path(), "attacker"); let (server, _index) = spawn_server(registered_line, MaxWireBytes::new(1 << 30)).await; let url = format!( "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", server.port );
let work = scratch.path().join("work"); let head = seed_work(&work); let (ok, out) = push(&work, &url, &attacker_path, &["main"]).await; assert!(!ok, "unauthorized push must be rejected:\n{out}");
let repo = server.layout.open(&server.repo_did).unwrap(); assert!( repo.references().unwrap().is_empty(), "denied push must create no ref" ); assert!( !repo.contains(Oid::from_hex(&head).unwrap()), "denied push must migrate no objects into the live odb" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn ref_namespace_policy() { let fx = fixture().await; seed_work(&fx.work);
for (spec, names_itself) in [ ("main:refs/hidden/feature/main", "server-side fork staging"), ("main:refs/atproto/commit", "atproto commit chain"), ( "main:refs/cob-checkpoints/sh.tangled.knot.member/beef", "knot-written collaborative-object snapshots", ), ] { let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &[spec]).await; assert!( !ok, "the ssh path screens namespaces too, so {spec} must be refused:\n{out}" ); assert!( out.contains(names_itself), "git must show the pusher which namespace refused {spec}:\n{out}" ); assert!( ref_names(&fx.server).is_empty(), "a push refused by policy must leave the ref store empty" ); }
let (ok, out) = push( &fx.work, &fx.url, &fx.key_path, &["main:refs/notes/commits"], ) .await; assert!( ok, "push to any non-reserved namespace must be accepted:\n{out}" ); assert!( ref_names(&fx.server) .iter() .any(|name| name == "refs/notes/commits"), "pushed ref must land" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn cob_ref_guard_lifecycle() { let fx = fixture().await; seed_work(&fx.work); let home = CobHome::from(&RepoDid::new(REPO_DID).unwrap()); let foreign = CobHome::from(&RepoDid::new("did:plc:whelk").unwrap());
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "head must land for the advertisement check:\n{out}");
let (owned_tip, owned_ref, owned_spec) = seed_cob(&fx.work, 1, "did:plc:limpet", &home); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &[owned_spec.as_str()]).await; assert!( ok, "COB ref signed by the repository key must verify and land over ssh:\n{out}" );
let (_forged_tip, forged_ref, forged_spec) = seed_cob(&fx.work, 9, "did:plc:whelk", &home); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &[forged_spec.as_str()]).await; assert!( !ok, "COB ref signed by a stranger must be refused at the receive boundary:\n{out}" );
let (_transplant_tip, transplant_ref, transplant_spec) = seed_cob(&fx.work, 1, "did:plc:mussel", &foreign); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &[transplant_spec.as_str()]).await; assert!( !ok, "same key signing for another repo's home must be refused on transplant:\n{out}" );
let landed = ref_names(&fx.server); assert!( landed.contains(&owned_ref), "owner-signed COB ref must be stored: {landed:?}" ); assert!( !landed.contains(&forged_ref), "stranger-signed COB ref must be absent: {landed:?}" ); assert!( !landed.contains(&transplant_ref), "transplanted COB ref must be absent: {landed:?}" );
let cob_name = RefName::new(&owned_ref).unwrap(); let del = format!(":{owned_ref}"); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &[del.as_str()]).await; assert!(!ok, "deleting a COB ref must be refused:\n{out}"); assert!( out.contains("append-only"), "rejection must name the append-only rule:\n{out}" ); assert!( ref_names(&fx.server).contains(&owned_ref), "COB ref must survive the refused delete" );
let repo = Repo::open(&fx.work).unwrap(); CobStore::new(&repo) .update( &home, knot_types::CobId::new(owned_tip), &MembersChange::Add(Grant { subject: AccountDid::new("did:plc:bailey").unwrap(), added_by: AccountDid::new(OWNER_DID).unwrap(), created_at: UnixSeconds::new(2), }), &K256Signer::generate(&SeededEntropy::new(1)), UnixSeconds::new(2), ) .unwrap(); assert_ne!( repo.find_ref(&cob_name).unwrap(), Some(owned_tip), "local COB ref now points at a new, equally valid tip" ); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &[owned_spec.as_str()]).await; assert!( !ok, "re-pushing a moved COB ref must be refused instead of silently clobbered:\n{out}" ); assert_eq!( fx.server .layout .open(&fx.server.repo_did) .unwrap() .find_ref(&cob_name) .unwrap(), Some(owned_tip), "live COB ref must still point at the original tip" );
let (ok, advert) = git_ssh(Path::new("/tmp"), &fx.key_path, &["ls-remote", &fx.url]).await; assert!(ok, "ls-remote over ssh must succeed:\n{advert}"); assert!( advert.contains("refs/heads/main"), "head must be advertised:\n{advert}" ); assert!( !advert.contains("refs/cobs/"), "no refs/cobs/* may leak into the ssh advertisement:\n{advert}" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_filled_key_set_refuses_an_unregistered_key_and_an_acl_write_reopens_the_check() { let fx = fixture().await; let head = seed_work(&fx.work); let head_oid = Oid::from_hex(&head).unwrap();
fx.index.keys().mark_ready(fx.index.generation()); fx.index.refresh_members().unwrap(); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!( ok, "a grant written after the key set was read must reopen the check, or whoever it \ grants is refused at the handshake until the next fill pass:\n{out}" ); assert_eq!( main_tip(&fx.server.layout, &fx.server.repo_did), Some(head_oid) );
fx.index.keys().record( &AccountDid::new(OWNER_DID).unwrap(), vec![knot_types::OfferedKey::from_bytes(registered_blob(&fx))], forever(), ); fx.index.keys().mark_ready(fx.index.generation());
let (unregistered_path, _unregistered_line) = keygen(fx.scratch.path(), "unregistered"); let two_ids = format!( "ssh -i {unregistered_path} -i {} {UNSHARED_SSH}", fx.key_path ); let (ok, out) = { let (work, url) = (fx.work.clone(), fx.url.clone()); tokio::task::spawn_blocking(move || { git( &work, &[("GIT_SSH_COMMAND", &two_ids)], &["push", "-q", &url, "main:refs/heads/second"], ) }) .await .unwrap() }; assert!( ok, "a filled key set refuses the unregistered key, so the client offers its registered key \ without the url identifying anybody:\n{out}" ); assert_eq!( fx.server .layout .open(&fx.server.repo_did) .unwrap() .find_ref(&RefName::new("refs/heads/second").unwrap()) .unwrap(), Some(head_oid) );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn key_recognition_edge_cases() { let fx = fixture().await; let head = seed_work(&fx.work); let head_oid = Oid::from_hex(&head).unwrap();
let (unregistered_path, _unregistered_line) = keygen(fx.scratch.path(), "unregistered"); let two_ids = format!( "ssh -i {unregistered_path} -i {} {UNSHARED_SSH}", fx.key_path ); let (ok, out) = { let (work, url) = (fx.work.clone(), fx.url.clone()); tokio::task::spawn_blocking(move || { git( &work, &[("GIT_SSH_COMMAND", &two_ids)], &["push", "-q", &url, "main"], ) }) .await .unwrap() }; assert!( !ok, "with the key set still filling, the push is checked against whichever key the client \ offers first:\n{out}" ); assert!( out.contains("@nel.pet"), "refusal lists who may push, so the pusher knows which key to offer:\n{out}" ); assert!( out.contains("IdentitiesOnly"), "refusal states how a multi-key client can offer its registered key:\n{out}" ); assert_eq!( main_tip(&fx.server.layout, &fx.server.repo_did), None, "the refused push leaves the repo empty" );
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "offering only the registered key must succeed:\n{out}"); assert_eq!( main_tip(&fx.server.layout, &fx.server.repo_did), Some(head_oid) );
let blob = russh::keys::ssh_key::PublicKey::from_openssh( &std::fs::read_to_string(fx.scratch.path().join("client.pub")).unwrap(), ) .unwrap() .to_bytes() .unwrap(); fx.index.keys().record( &AccountDid::new("did:plc:cuttle").unwrap(), vec![knot_types::OfferedKey::from_bytes(blob)], forever(), ); let (ok, out) = push( &fx.work, &fx.url, &fx.key_path, &["main:refs/heads/squat-check"], ) .await; assert!( ok, "stranger who published the owner's key mustn't deny the owner's push:\n{out}" ); assert_eq!( fx.server .layout .open(&fx.server.repo_did) .unwrap() .find_ref(&RefName::new("refs/heads/squat-check").unwrap()) .unwrap(), Some(head_oid) );}
#[test]fn a_group_or_other_readable_host_key_is_refused_on_load() { use std::os::unix::fs::PermissionsExt; let dir = tempfile::tempdir().unwrap(); let path = dir.path().join("host"); knot_ssh::load_or_create_host_key(&path).unwrap(); assert_eq!( std::fs::metadata(&path).unwrap().permissions().mode() & 0o777, 0o600, "freshly created host key is 0600" );
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o644)).unwrap(); let refused = knot_ssh::load_or_create_host_key(&path); assert!( matches!(refused, Err(knot_ssh::SshError::HostKey { .. })), "world-readable existing host key must be refused on load: {refused:?}" );}
async fn launch(host_key_dir: &Path, layout: Layout, index: Arc<Index>, accounts: Accounts) -> u16 { let atproto = Arc::new(Atproto::new( multi_http(accounts), ManualClock::new(UnixMicros::new(1_000_000_000)), KnotId::new("did:web:nel.pet").unwrap(), knot_atproto::PlcDirectory::new(Url::parse("https://plc.directory/").unwrap()).unwrap(), )); std::fs::create_dir_all(host_key_dir).unwrap(); let host_key = knot_ssh::load_or_create_host_key(&host_key_dir.join("host")).unwrap(); let state = Arc::new(knot_ssh::SshState::new(knot_ssh::SshConfig { layout, index, atproto, knot_actor: actor_for_seed(77), hostname: knot_types::KnotHostname::new("knot.test").unwrap(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), admins: std::collections::BTreeSet::new(), admission: knot_types::AdmissionPolicy::Closed, max_pack_bytes: MaxWireBytes::new(1 << 30), archive_limit: ArchiveLimit::default(), ci_logs: None, })); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let port = listener.local_addr().unwrap().port(); tokio::spawn(async move { let _ = knot_ssh::serve_on_socket(listener, host_key, state).await; }); port}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_handle_in_the_url_identifies_a_visitor_and_lets_a_multi_key_client_find_its_key() { let fx = fixture().await; let head = seed_work(&fx.work); let head_oid = Oid::from_hex(&head).unwrap();
let port = fx.server.port; let greeted_key = fx.key_path.clone(); let (_ok, out) = tokio::task::spawn_blocking(move || ssh_bare_as(&greeted_key, "nel.pet", port)) .await .unwrap(); assert!( out.contains("@nel.pet"), "an asserted handle identifies the visitor on first contact, with an empty cache:\n{out}" );
let (unregistered_path, _unregistered_line) = keygen(fx.scratch.path(), "unregistered"); let two_ids = format!( "ssh -i {unregistered_path} -i {} {UNSHARED_SSH}", fx.key_path ); let identified = format!("ssh://nel.pet@127.0.0.1:{}/{REPO_DID}", fx.server.port); let (ok, out) = { let (work, url) = (fx.work.clone(), identified.clone()); tokio::task::spawn_blocking(move || { git( &work, &[("GIT_SSH_COMMAND", &two_ids)], &["push", "-q", &url, "main"], ) }) .await .unwrap() }; assert!( ok, "a handle in the url lets the knot refuse the unregistered key so the client offers the \ next key:\n{out}" ); assert_eq!( main_tip(&fx.server.layout, &fx.server.repo_did), Some(head_oid) ); assert_eq!( fx.index.owner_of_key( &knot_types::OfferedKey::from_bytes(registered_blob(&fx)), UnixSeconds::new(0), ), Resolved::Ready(None), "an asserted handle is whatever the client typed, so the keys read for it mustn't enter \ the set, or anyone can fill the key budget by asserting handles" );}
fn registered_blob(fx: &Fixture) -> Vec<u8> { russh::keys::ssh_key::PublicKey::from_openssh( &std::fs::read_to_string(fx.scratch.path().join("client.pub")).unwrap(), ) .unwrap() .to_bytes() .unwrap()}
fn registered_index( scratch: &TempDir, budget: knot_index::KeyBudget,) -> (Layout, RepoDid, Arc<Index>) { let meta_path = scratch.path().join("meta"); Repo::create(&meta_path).unwrap(); let layout = Layout::new(scratch.path().join("repos")); let repo_did = RepoDid::new(REPO_DID).unwrap(); layout.create(&repo_did).unwrap();
let signer = K256Signer::generate(&SeededEntropy::new(2)); let meta = Repo::open(&meta_path).unwrap(); CobStore::new(&meta) .create( &CobHome::from(&KnotId::new("did:web:nel.pet").unwrap()), &RegistryChange::Register(Registration { owner: OwnerDid::new(OWNER_DID).unwrap(), rkey: RepoRkey::new(REPO_NAME).unwrap(), name: RepoName::new(REPO_NAME).unwrap(), repo: repo_did.clone(), created_at: UnixSeconds::new(1), policy: knot_types::ContributionPolicy::default(), }), &signer, UnixSeconds::new(1), ) .unwrap();
let index = Arc::new(Index::with_key_budget(meta_path, layout.clone(), budget)); index.rebuild().unwrap(); index.warm_collaborators(); (layout, repo_did, index)}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn an_unreachable_pds_reads_as_transient_and_its_keys_get_in_once_it_recovers() { let scratch = tempfile::tempdir().unwrap(); let (stale_key, stale_line) = keygen(scratch.path(), "stale"); let (fresh_key, fresh_line) = keygen(scratch.path(), "fresh"); let (layout, repo_did, index) = registered_index(&scratch, knot_index::KeyBudget::DEFAULT);
let accounts = Accounts::publishing(HashMap::from([(OWNER_DID.to_string(), vec![stale_line])])) .unreachable(HashSet::from([OWNER_DID.to_string()])); let port = launch( &scratch.path().join("hostkey"), layout.clone(), Arc::clone(&index), accounts.clone(), ) .await; let url = format!("ssh://git@127.0.0.1:{port}/{REPO_DID}");
let work = scratch.path().join("work"); let head = seed_work(&work); let (ok, out) = push(&work, &url, &stale_key, &["main"]).await; assert!( !ok, "a push mustn't be accepted while the owner's records are unreadable:\n{out}" ); assert!( out.contains("retry shortly"), "an unreadable PDS must read as transient:\n{out}" ); assert!( !out.contains("doesn't match"), "a transient failure mustn't be reported to the pusher as a wrong key:\n{out}" ); assert_eq!( main_tip(&layout, &repo_did), None, "the refused push leaves the repo empty" );
accounts.restore(OWNER_DID); let (ok, out) = push(&work, &url, &stale_key, &["main"]).await; assert!( ok, "the key the owner publishes must push once its PDS answers again:\n{out}" );
accounts.publish(OWNER_DID, fresh_line); let (ok, out) = push(&work, &url, &fresh_key, &["main", "--force"]).await; assert!( ok, "a key the owner published after the knot last read the account must get in on the next \ push, or publishing a second key locks its owner out until a fill pass catches up:\n{out}" ); assert_eq!( main_tip(&layout, &repo_did), Some(Oid::from_hex(&head).unwrap()) );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn an_account_the_budget_couldnt_fit_still_clears_the_handshake_and_pushes() { let scratch = tempfile::tempdir().unwrap(); let (owner_key, owner_line) = keygen(scratch.path(), "owner"); let (layout, repo_did, index) = registered_index(&scratch, knot_index::KeyBudget::from_bytes(200));
let blob = russh::keys::ssh_key::PublicKey::from_openssh(&owner_line) .unwrap() .to_bytes() .unwrap(); assert_eq!( index.keys().record( &AccountDid::new(OWNER_DID).unwrap(), vec![knot_types::OfferedKey::from_bytes(blob)], forever(), ), knot_index::KeyRecord::Unheld, "a 200-byte budget records the read without keeping the key" ); index.keys().mark_ready(index.generation());
let port = launch( &scratch.path().join("hostkey"), layout.clone(), Arc::clone(&index), Accounts::publishing(HashMap::from([(OWNER_DID.to_string(), vec![owner_line])])), ) .await; let url = format!("ssh://git@127.0.0.1:{port}/{REPO_DID}");
let work = scratch.path().join("work"); let head = seed_work(&work); let (ok, out) = push(&work, &url, &owner_key, &["main"]).await; assert!( ok, "the set can't fit the owner's keys, so the handshake must defer to the push check \ instead of refusing a key the accounts on file don't publish:\n{out}" ); assert_eq!( main_tip(&layout, &repo_did), Some(Oid::from_hex(&head).unwrap()) );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_collaborator_pushes_its_repo_but_is_denied_on_a_repo_it_doesnt_collaborate_on() { const REPO_A: &str = "did:plc:squid"; const REPO_B: &str = "did:plc:clam"; const OWNER: &str = "did:plc:nel"; const COLLAB: &str = "did:plc:olaren";
let scratch = tempfile::tempdir().unwrap(); let (owner_key, owner_line) = keygen(scratch.path(), "owner"); let (collab_key, collab_line) = keygen(scratch.path(), "collab");
let meta_path = scratch.path().join("meta"); Repo::create(&meta_path).unwrap(); let layout = Layout::new(scratch.path().join("repos")); let repo_a = RepoDid::new(REPO_A).unwrap(); let repo_b = RepoDid::new(REPO_B).unwrap(); let git_a = layout.create(&repo_a).unwrap(); layout.create(&repo_b).unwrap();
let signer = K256Signer::generate(&SeededEntropy::new(2)); let meta = Repo::open(&meta_path).unwrap(); let store = CobStore::new(&meta); let knot_home = CobHome::from(&KnotId::new("did:web:nel.pet").unwrap()); let reg = store .create( &knot_home, &RegistryChange::Register(Registration { owner: OwnerDid::new(OWNER).unwrap(), rkey: RepoRkey::new("anemone").unwrap(), name: RepoName::new("anemone").unwrap(), repo: repo_a.clone(), created_at: UnixSeconds::new(1), policy: knot_types::ContributionPolicy::default(), }), &signer, UnixSeconds::new(1), ) .unwrap(); store .update( &knot_home, reg.object, &RegistryChange::Register(Registration { owner: OwnerDid::new(OWNER).unwrap(), rkey: RepoRkey::new("barnacle").unwrap(), name: RepoName::new("barnacle").unwrap(), repo: repo_b.clone(), created_at: UnixSeconds::new(2), policy: knot_types::ContributionPolicy::default(), }), &signer, UnixSeconds::new(2), ) .unwrap(); store .create( &knot_home, &MembersChange::Add(Grant { subject: AccountDid::new(COLLAB).unwrap(), added_by: AccountDid::new(OWNER).unwrap(), created_at: UnixSeconds::new(1), }), &signer, UnixSeconds::new(1), ) .unwrap(); CobStore::new(&git_a) .create( &CobHome::from(&repo_a), &CollaboratorsChange::Add(Grant { subject: AccountDid::new(COLLAB).unwrap(), added_by: AccountDid::new(OWNER).unwrap(), created_at: UnixSeconds::new(1), }), &signer, UnixSeconds::new(1), ) .unwrap();
let index = Arc::new(Index::new(meta_path, layout.clone())); index.rebuild().unwrap(); index.warm_collaborators();
let identities = HashMap::from([ (OWNER.to_string(), vec![owner_line]), (COLLAB.to_string(), vec![collab_line]), ]); let accounts = Accounts::publishing(identities); let port = launch( &scratch.path().join("hostkey"), layout.clone(), Arc::clone(&index), accounts.clone(), ) .await;
let work_a = scratch.path().join("work_a"); let head_a = seed_work(&work_a); let url_a = format!("ssh://git@127.0.0.1:{port}/{REPO_A}"); let (ok, out) = push(&work_a, &url_a, &collab_key, &["main"]).await; assert!( ok, "collaborator must push the repo it collaborates on:\n{out}" ); assert_eq!( main_tip(&layout, &repo_a), Some(Oid::from_hex(&head_a).unwrap()), "collaborator's commit must be repo A's main tip" );
let work_b = scratch.path().join("work_b"); seed_work(&work_b); let url_b = format!("ssh://git@127.0.0.1:{port}/{REPO_B}"); let (denied, out) = push(&work_b, &url_b, &collab_key, &["main"]).await; assert!( !denied, "a key that pushes repo A must be denied on repo B, where its owner was never granted:\n{out}" ); assert!( main_tip(&layout, &repo_b).is_none(), "denied cross-repo push must land nothing on repo B" );
let after_first_denial = accounts.listings(); let (denied, out) = push(&work_b, &url_b, &collab_key, &["main"]).await; assert!(!denied, "the second attempt is denied the same way:\n{out}"); assert_eq!( accounts.listings(), after_first_denial, "repo B's owner was read during the first denial and is on file, so retrying mustn't \ read that PDS again, or anyone with a key can make the knot fetch from a third party \ at will:\n{out}" );
let work_owner = scratch.path().join("work_owner_b"); let head_owner = seed_work(&work_owner); let (ok, out) = push(&work_owner, &url_b, &owner_key, &["main"]).await; assert!(ok, "owner must push to repo B:\n{out}"); assert_eq!( main_tip(&layout, &repo_b), Some(Oid::from_hex(&head_owner).unwrap()), "owner's push to repo B must land, isolating the collaborator's denial as authorization" );}
const NEW_COLLAB: &str = "did:plc:olaren";const COLLAB_BYSTANDER: &str = "did:plc:bailey";
struct Invitation { scratch: TempDir, layout: Layout, url: String, head: Oid, collab_key: String, bystander_key: String, accounts: Accounts,}
impl Invitation { fn tip(&self) -> Option<Oid> { main_tip(&self.layout, &RepoDid::new(REPO_DID).unwrap()) } async fn push(&self, key: &str) -> (bool, String) { let work = self.scratch.path().join("work"); push(&work, &self.url, key, &["main"]).await }}
async fn a_new_collaborator(pds_down: bool) -> Invitation { let scratch = tempfile::tempdir().unwrap(); let dir = scratch.path(); let (_owner_key, owner_line) = keygen(dir, "owner"); let (collab_key, collab_line) = keygen(dir, "collab"); let (bystander_key, bystander_line) = keygen(dir, "bystander"); let (layout, repo, index) = registered_index(&scratch, knot_index::KeyBudget::DEFAULT); let head = Oid::from_hex(&seed_work(&dir.join("work"))).unwrap(); let now = UnixSeconds::new(1); let invited = CollaboratorsChange::Add(Grant { subject: AccountDid::new(NEW_COLLAB).unwrap(), added_by: AccountDid::new(OWNER_DID).unwrap(), created_at: now, }); let signer = K256Signer::generate(&SeededEntropy::new(2)); CobStore::new(&layout.open(&repo).unwrap()) .create(&CobHome::from(&repo), &invited, &signer, now) .unwrap(); index.refresh_collaborators(&repo).unwrap();
let held = |line: &str| { let key = russh::keys::ssh_key::PublicKey::from_openssh(line).unwrap(); vec![knot_types::OfferedKey::from_bytes(key.to_bytes().unwrap())] }; let owner = AccountDid::new(OWNER_DID).unwrap(); let bystander = AccountDid::new(COLLAB_BYSTANDER).unwrap(); let keys = index.keys(); keys.record(&owner, held(&owner_line), forever()); keys.record(&bystander, held(&bystander_line), forever());
let records = Accounts::publishing(HashMap::from([ (OWNER_DID.to_string(), vec![owner_line]), (NEW_COLLAB.to_string(), vec![collab_line]), ])); let accounts = match pds_down { false => records, true => records.unreachable(HashSet::from([NEW_COLLAB.to_string()])), }; let port = launch(dir, layout.clone(), Arc::clone(&index), accounts.clone()).await; Invitation { url: format!("ssh://git@127.0.0.1:{port}/{REPO_DID}"), scratch, layout, head, collab_key, bystander_key, accounts, }}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_new_collaborator_pushes_once_the_knot_reads_their_keys() { let invitation = a_new_collaborator(false).await; let (ok, out) = invitation.push(&invitation.collab_key).await; assert!( ok, "the roll seats this account and the push is what makes the knot read their keys, so \ no second call from the client is needed:\n{out}" ); assert_eq!(invitation.tip(), Some(invitation.head));}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn wide_door_admits_whichever_key_a_new_collaborator_offers_first() { let invitation = a_new_collaborator(false).await; let (decoy, _line) = keygen(invitation.scratch.path(), "decoy"); let offered = format!("{decoy} -i {}", invitation.collab_key); let (ok, out) = invitation.push(&offered).await; assert!( !ok, "the ssh username is git and the repository arrives in the exec command, so at key \ identification the knot can't know which roll to weigh a key against:\n{out}" ); assert!( out.contains("IdentitiesOnly"), "accepting the offer ends the client's key search, so the refusal has to tell a \ collaborator holding two keys how to offer their registered one:\n{out}" ); assert!( invitation.tip().is_none(), "main moved, so the decoy key pushed" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_new_collaborators_first_push_survives_a_strangers_failed_one() { let invitation = a_new_collaborator(true).await; let (ok, out) = invitation.push(&invitation.bystander_key).await; assert!( !ok, "bystander's push is refused, and roll doesn't list their key anywhere" ); assert!( out.contains("couldn't read the account records"), "key is on file: session doesn't cost a miss reservation, and fan-out runs, which is what spends the collaborator's pacer while their pds is down:\n{out}" );
invitation.accounts.restore(NEW_COLLAB); let (ok, out) = invitation.push(&invitation.collab_key).await; assert!( ok, "a collaborator with no keys on file is read on every attempt, so a pacer somebody else \ spent can't turn their first push into an unregistered key:\n{out}" ); assert_eq!(invitation.tip(), Some(invitation.head));}
fn ssh_bare(key_path: &str, port: u16) -> (bool, String) { ssh_bare_as(key_path, "git", port)}
fn ssh_bare_as(key_path: &str, user: &str, port: u16) -> (bool, String) { let out = Command::new("ssh") .args(["-i", key_path]) .args(UNSHARED_SSH.split_whitespace()) .args(["-p", &port.to_string(), &format!("{user}@127.0.0.1")]) .output() .expect("ssh runs"); ( out.status.success(), format!( "{}{}", String::from_utf8_lossy(&out.stdout), String::from_utf8_lossy(&out.stderr) ), )}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_bare_ssh_session_greets_a_visitor_then_identifies_them_once_they_have_pushed() { let fx = fixture().await; let port = fx.server.port;
let key_path = fx.key_path.clone(); let (_ok, out) = tokio::task::spawn_blocking(move || ssh_bare(&key_path, port)) .await .unwrap(); assert!(out.contains("knot.test"), "greeting names the knot:\n{out}"); assert!( out.contains("ssh key"), "a visitor the knot can't identify yet learns what a push needs:\n{out}" );
seed_work(&fx.work); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "seeding main must succeed:\n{out}");
let key_path = fx.key_path.clone(); let (_ok, out) = tokio::task::spawn_blocking(move || ssh_bare(&key_path, port)) .await .unwrap(); assert!( out.contains("@nel.pet"), "a push teaches the knot the key, so the next greeting uses the handle:\n{out}" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_push_to_a_new_branch_offers_a_pull_request_link() { let fx = fixture().await; seed_work(&fx.work);
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "seeding main must land:\n{out}");
git(&fx.work, &[], &["checkout", "-q", "-b", "feature"]); std::fs::write(fx.work.join("feature.txt"), "work\n").unwrap(); git(&fx.work, &[], &["add", "-A"]); git(&fx.work, &[], &["commit", "-q", "-m", "feature work"]);
let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["feature"]).await; assert!(ok, "feature-branch push must land:\n{out}"); assert!( out.contains("https://tangled.test/nel.pet/anemone/pulls/new"), "new non-default branch is answered with a pull-request link:\n{out}" ); assert!( out.contains("sourceBranch=feature") && out.contains("targetBranch=main"), "link points the new branch at the default:\n{out}" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn a_verbose_ci_push_option_reports_a_clean_pipeline() { let fx = fixture().await; std::fs::create_dir_all(fx.work.join(".tangled/workflows")).unwrap(); git(&fx.work, &[], &["init", "-q", "-b", "main"]); std::fs::write( fx.work.join(".tangled/workflows/ci.yml"), "engine: nixery.dev/x\nwhen:\n - event: push\n branch: ['**']\n", ) .unwrap(); git(&fx.work, &[], &["add", "-A"]); git(&fx.work, &[], &["commit", "-q", "-m", "add ci"]);
let (ok, out) = push( &fx.work, &fx.url, &fx.key_path, &["--push-option=verbose-ci", "main"], ) .await; assert!(ok, "push with a push option must land:\n{out}"); assert!( out.contains("no diagnostics"), "verbose-ci reports clean compile over the sideband:\n{out}" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn git_archive_remote_over_ssh_streams_a_tar_of_the_tree() { let fx = fixture().await; seed_work(&fx.work); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "the seeding push must succeed before archiving:\n{out}");
let out_tar = fx.scratch.path().join("archive.tar"); let (ok, out) = git_ssh( &fx.work, &fx.key_path, &[ "archive", "--format=tar", "--remote", &fx.url, "-o", out_tar.to_str().unwrap(), "HEAD", ], ) .await; assert!(ok, "git archive --remote over ssh must succeed:\n{out}");
let tar = std::fs::read(&out_tar).unwrap(); assert!( knot_fixtures::contains(&tar, b"README.md"), "archived tar must contain the README.md entry" );}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn git_archive_remote_over_ssh_honors_the_configured_archive_limit() { let fx = fixture_with_archive_limit(ArchiveLimit::new(512)).await; seed_work(&fx.work); let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; assert!(ok, "the seeding push must succeed before archiving:\n{out}");
let out_tar = fx.scratch.path().join("archive.tar"); let (ok, out) = git_ssh( &fx.work, &fx.key_path, &[ "archive", "--format=tar", "--remote", &fx.url, "-o", out_tar.to_str().unwrap(), "HEAD", ], ) .await; assert!(!ok, "git archive --remote past the limit must fail:\n{out}"); assert!( out.contains("archive exceeds the 512 byte limit"), "the refusal must reach the client over the ssh channel:\n{out}" );}
fn pkt(payload: &[u8]) -> Vec<u8> { let mut framed = format!("{:04x}", payload.len() + 4).into_bytes(); framed.extend_from_slice(payload); framed}
fn pkt_text(line: &str) -> Vec<u8> { pkt(format!("{line}\n").as_bytes())}
fn read_until(reader: &mut impl std::io::Read, needle: &[u8], buffer: &mut Vec<u8>) { std::iter::from_fn(|| { let mut byte = [0u8; 1]; match reader.read(&mut byte) { Ok(0) | Err(_) => None, Ok(_) => { buffer.push(byte[0]); Some(buffer.ends_with(needle)) } } }) .find(|done| *done) .expect("the session must answer before closing the stream");}
fn trickled_lfs_upload( key_path: &str, port: u16, body: &[u8], oid: &str, midway: std::sync::mpsc::Sender<()>,) -> (bool, String) { use std::io::Write; let mut child = Command::new("ssh") .args(["-i", key_path]) .args(UNSHARED_SSH.split_whitespace()) .args([ "-p", &port.to_string(), "git@127.0.0.1", &format!("git-lfs-transfer '{OWNER_DID}/{REPO_NAME}' upload"), ]) .stdin(std::process::Stdio::piped()) .stdout(std::process::Stdio::piped()) .stderr(std::process::Stdio::null()) .spawn() .expect("ssh runs"); let mut stdin = child.stdin.take().unwrap(); let mut stdout = child.stdout.take().unwrap(); let mut transcript = Vec::new();
read_until(&mut stdout, b"version=1\n0000", &mut transcript);
let (first, second) = body.split_at(body.len() / 2); stdin .write_all(&pkt_text(&format!("put-object {oid}"))) .unwrap(); stdin .write_all(&pkt_text(&format!("size={}", body.len()))) .unwrap(); stdin.write_all(b"0001").unwrap(); first.chunks(32 * 1024).for_each(|chunk| { stdin.write_all(&pkt(chunk)).unwrap(); }); stdin.flush().unwrap(); midway.send(()).unwrap(); std::thread::sleep(std::time::Duration::from_millis(900));
second.chunks(32 * 1024).for_each(|chunk| { stdin.write_all(&pkt(chunk)).unwrap(); }); stdin.write_all(b"0000").unwrap(); stdin.flush().unwrap(); read_until(&mut stdout, b"status 200\n0000", &mut transcript);
stdin.write_all(&pkt_text("quit")).unwrap(); stdin.write_all(b"0000").unwrap(); stdin.flush().unwrap(); drop(stdin); use std::io::Read; let _ = stdout.read_to_end(&mut transcript); let status = child.wait().expect("ssh exits"); ( status.success(), String::from_utf8_lossy(&transcript).into_owned(), )}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn shutdown_drains_an_in_flight_lfs_transfer_before_exit() { use knot_lfs::{LfsStore, Pool}; let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "drain"); let lfs_dir = scratch.path().join("lfs"); std::fs::create_dir_all(&lfs_dir).unwrap(); let handle = knot_lfs::LfsHandle::open( knot_lfs::LfsStorePath::new(&lfs_dir), knot_lfs::LfsSize::new(1 << 30), knot_lfs::FreeSpaceFloor::new(0), ) .unwrap(); let (server, _index, shutdown, serve_task, _) = spawn_server_core( public_line, MaxWireBytes::new(1 << 20), ArchiveLimit::default(), true, Some(handle.clone()), false, ) .await;
let body: Vec<u8> = (0..1_048_576u32).map(|n| (n % 251) as u8).collect(); let oid = knot_lfs::LfsOid::from_digest(knot_types::Sha256Digest::hash(&body)); let (midway_tx, midway_rx) = std::sync::mpsc::channel();
let client = { let key_path = key_path.clone(); let oid = oid.clone(); let port = server.port; tokio::task::spawn_blocking(move || { trickled_lfs_upload(&key_path, port, &body, oid.as_str(), midway_tx) }) };
tokio::task::spawn_blocking(move || { midway_rx .recv_timeout(std::time::Duration::from_secs(20)) .expect("the upload must reach its midway point") }) .await .unwrap();
shutdown.cancel(); tokio::time::sleep(std::time::Duration::from_millis(150)).await; assert!( !serve_task.is_finished(), "the listener must keep draining while a transfer is in flight" );
let (ok, transcript) = client.await.unwrap(); assert!( ok, "the in-flight upload must finish cleanly across the shutdown:\n{transcript}" ); assert!( transcript.contains("status 200"), "the server must acknowledge the drained upload:\n{transcript}" );
tokio::time::timeout(std::time::Duration::from_secs(10), serve_task) .await .expect("the drained listener must exit promptly once transfers finish") .unwrap();
let repo_did = RepoDid::new(REPO_DID).unwrap(); assert_eq!( handle .store .probe(Pool::Lfs, &repo_did, &oid) .unwrap() .map(|size| size.get()), Some(1_048_576), "the drained upload must be durable" );
let (connected, _) = { let key_path = key_path.clone(); let port = server.port; tokio::task::spawn_blocking(move || ssh_bare(&key_path, port)) .await .unwrap() }; assert!( !connected, "a connection after shutdown must be refused, the drain only covers in-flight work" );}
#[derive(serde::Deserialize)]struct OpProbe { path: String,}
#[derive(serde::Deserialize)]struct CommitProbe { ops: Vec<OpProbe>,}
fn frame_kind_and_ops(frame: &[u8]) -> (&'static str, Vec<String>) { let mut read = frame; let header: serde_json::Value = serde_ipld_dagcbor::de::from_reader_once(&mut read).unwrap(); match header["t"].as_str().unwrap_or_default() { "#commit" => { let payload: CommitProbe = serde_ipld_dagcbor::de::from_reader_once(&mut read).unwrap(); let ops = payload.ops.into_iter().map(|op| op.path).collect(); ("#commit", ops) } "#sync" => ("#sync", Vec::new()), _ => ("#ignored", Vec::new()), }}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]async fn an_authorized_ssh_push_projects_the_ref_records_the_http_push_projects() { let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "client"); let (server, _index, _shutdown, _serve_task, firehose) = spawn_server_core( public_line, MaxWireBytes::new(1 << 30), ArchiveLimit::default(), true, None, true, ) .await; let firehose = firehose.expect("the projection fixture serves a firehose");
let url = format!( "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", server.port ); let work = scratch.path().join("work"); let head = seed_work(&work);
let (ok, out) = push(&work, &url, &key_path, &["main"]).await; assert!(ok, "authorized ssh push must succeed:\n{out}"); assert_eq!( main_tip(&server.layout, &server.repo_did), Some(Oid::from_hex(&head).unwrap()), "pushed commit must be the repository's main tip" );
let ref_ops: Vec<String> = firehose .replay( knot_events::EventCursor::START, knot_events::ReplayBounds::new( knot_events::ReplayEvents::new(4096).unwrap(), knot_events::ReplayBytes::new(16 << 20).unwrap(), ), ) .events .iter() .flat_map(|event| frame_kind_and_ops(&event.frame).1) .filter(|path| path.starts_with("sh.tangled.git.ref/")) .collect(); assert_eq!( ref_ops, vec!["sh.tangled.git.ref/refs~2fheads~2fmain".to_string()], "the ssh push projected its ref into an atproto record, exactly as the http push does" );}