This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-ssh / examples / ephemeral_knot.rs
8.1 kB 243 lines
1use std::collections::BTreeSet; 2use std::sync::Arc; 3use std::time::Duration; 4 5use knot_atproto::Atproto; 6use knot_cob::{CobHome, CobStore}; 7use knot_cobs::{Registration, RegistryChange}; 8use knot_git::{Layout, Repo}; 9use knot_index::Index; 10use knot_runtime::{ 11 FakeHttp, HttpResponse, K256Signer, ManualClock, SeededEntropy, Signer, UnixMicros, 12}; 13use knot_types::{ 14 AdmissionPolicy, KnotHostname, KnotId, OwnerDid, RepoDid, RepoName, RepoRkey, UnixSeconds, 15}; 16use tokio::net::TcpListener; 17use url::Url; 18 19#[global_allocator] 20static GLOBAL: tikv_jemallocator::Jemalloc = tikv_jemallocator::Jemalloc; 21 22const REPO_DID: &str = "did:plc:squid"; 23const REPO_NAME: &str = "anemone"; 24const OWNER_DID: &str = "did:plc:nel"; 25const PDS_HOST: &str = "pds.oyster.cafe"; 26 27fn apply_decay(ms: isize) { 28 unsafe { 29 let _ = tikv_jemalloc_ctl::raw::write(b"arenas.dirty_decay_ms\0", ms); 30 if let Ok(narenas) = tikv_jemalloc_ctl::raw::read::<u32>(b"arenas.narenas\0") { 31 (0..narenas).for_each(|arena| { 32 let name = format!("arena.{arena}.dirty_decay_ms\0"); 33 let _ = tikv_jemalloc_ctl::raw::write(name.as_bytes(), ms); 34 }); 35 } 36 } 37} 38 39async fn govern_decay() { 40 unsafe { 41 let _ = tikv_jemalloc_ctl::raw::write(b"background_thread\0", true); 42 } 43 let mut ticker = tokio::time::interval(Duration::from_secs(1)); 44 let mut applied = knot_resource::target_decay(); 45 apply_decay(applied.ms()); 46 loop { 47 ticker.tick().await; 48 let target = knot_resource::target_decay(); 49 if target != applied { 50 apply_decay(target.ms()); 51 applied = target; 52 } 53 } 54} 55 56fn vm_hwm_mib() -> u64 { 57 std::fs::read_to_string("/proc/self/status") 58 .unwrap_or_default() 59 .lines() 60 .find_map(|line| line.strip_prefix("VmHWM:")) 61 .and_then(|rest| rest.split_whitespace().next()) 62 .and_then(|kb| kb.parse::<u64>().ok()) 63 .map(|kb| kb / 1024) 64 .unwrap_or(0) 65} 66 67fn did_document(signer: &K256Signer, did: &str, pds: &str) -> Vec<u8> { 68 let multikey = knot_types::crypto::multikey(0xe7, signer.public_key().as_bytes()); 69 serde_json::to_vec(&serde_json::json!({ 70 "id": did, 71 "alsoKnownAs": ["at://nel.pet"], 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 87fn list_records_body(published_line: &str) -> Vec<u8> { 88 serde_json::to_vec(&serde_json::json!({ 89 "records": [{ 90 "value": { 91 "$type": "sh.tangled.publicKey", 92 "key": published_line, 93 "name": "laptop", 94 "createdAt": "2026-06-08T00:00:00Z" 95 } 96 }] 97 })) 98 .unwrap() 99} 100 101fn ok_body(body: Vec<u8>) -> HttpResponse { 102 HttpResponse { 103 status: http::StatusCode::OK, 104 headers: http::HeaderMap::new(), 105 body: bytes::Bytes::from(body), 106 } 107} 108 109fn not_found() -> HttpResponse { 110 HttpResponse { 111 status: http::StatusCode::NOT_FOUND, 112 headers: http::HeaderMap::new(), 113 body: bytes::Bytes::new(), 114 } 115} 116 117fn fake_http(published_line: String) -> impl knot_runtime::HttpTransport { 118 let signer = K256Signer::generate(&SeededEntropy::new(1)); 119 let pds = format!("https://{PDS_HOST}"); 120 FakeHttp::new(move |request| { 121 let host = request.url.host_str().unwrap_or_default().to_string(); 122 let path = request.url.path().to_string(); 123 let body = if host == PDS_HOST { 124 list_records_body(&published_line) 125 } else if path.ends_with(REPO_DID) { 126 did_document(&signer, REPO_DID, &pds) 127 } else if path.ends_with(OWNER_DID) { 128 did_document(&signer, OWNER_DID, &pds) 129 } else { 130 return Ok(not_found()); 131 }; 132 Ok(ok_body(body)) 133 }) 134} 135 136#[tokio::main(flavor = "multi_thread")] 137async fn main() { 138 let mut args = std::env::args().skip(1); 139 let port: u16 = args 140 .next() 141 .expect("usage: ephemeral_knot <port> <client-pubkey-file>") 142 .parse() 143 .expect("port"); 144 let pubkey_path = args.next().expect("client-pubkey-file"); 145 let published_line = std::fs::read_to_string(&pubkey_path) 146 .expect("read client pubkey") 147 .trim() 148 .to_string(); 149 150 let max_threads = std::env::var("KNOT_MAX_THREADS") 151 .ok() 152 .and_then(|value| value.parse::<usize>().ok()); 153 knot_resource::init(knot_resource::Ceilings { 154 max_threads: max_threads.map(knot_resource::ThreadCount::new), 155 max_memory: None, 156 }); 157 158 let scan = tempfile::tempdir().expect("scan tempdir"); 159 let meta_path = scan.path().join("meta"); 160 Repo::create(&meta_path).unwrap(); 161 let layout = Layout::new(scan.path().join("repos")); 162 let repo_did = RepoDid::new(REPO_DID).unwrap(); 163 layout.create(&repo_did).unwrap(); 164 165 let signer = K256Signer::generate(&SeededEntropy::new(2)); 166 let meta = Repo::open(&meta_path).unwrap(); 167 CobStore::new(&meta) 168 .create( 169 &CobHome::from(&KnotId::new("did:web:nel.pet").unwrap()), 170 &RegistryChange::Register(Registration { 171 owner: OwnerDid::new(OWNER_DID).unwrap(), 172 rkey: RepoRkey::new(REPO_NAME).unwrap(), 173 name: RepoName::new(REPO_NAME).unwrap(), 174 repo: repo_did.clone(), 175 created_at: UnixSeconds::new(1), 176 }), 177 &signer, 178 UnixSeconds::new(1), 179 ) 180 .unwrap(); 181 182 let index = Arc::new(Index::new(meta_path, layout.clone())); 183 index.rebuild().unwrap(); 184 185 let atproto = Arc::new(Atproto::new( 186 fake_http(published_line), 187 ManualClock::new(UnixMicros::new(1_000_000_000)), 188 KnotId::new("did:web:nel.pet").unwrap(), 189 knot_atproto::PlcDirectory::new(Url::parse("https://plc.directory/").unwrap()).unwrap(), 190 )); 191 192 let key_dir = scan.path().join("hostkey"); 193 std::fs::create_dir_all(&key_dir).unwrap(); 194 let host_key = knot_ssh::load_or_create_host_key(&key_dir.join("host")).unwrap(); 195 196 let events = Arc::new(knot_events::EventLog::new( 197 ManualClock::new(UnixMicros::new(1_000_000_000)), 198 knot_events::ReplayBounds::new( 199 knot_events::ReplayEvents::new(64).unwrap(), 200 knot_events::ReplayBytes::new(16 << 20).unwrap(), 201 ), 202 )); 203 let actor = knot_types::ActorId::from_secp256k1( 204 K256Signer::generate(&SeededEntropy::new(1)) 205 .public_key() 206 .as_bytes(), 207 ); 208 let state = Arc::new(knot_ssh::SshState::new(knot_ssh::SshConfig { 209 layout, 210 index, 211 atproto, 212 knot_actor: actor, 213 events, 214 hostname: KnotHostname::new("knot.test").unwrap(), 215 appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), 216 admins: BTreeSet::new(), 217 admission: AdmissionPolicy::Closed, 218 max_pack_bytes: knot_pack::MaxWireBytes::new(1 << 34), 219 archive_limit: knot_git::ArchiveLimit::default(), 220 languages_push_budget: knot_postreceive::LanguagesPushBudget::new(Duration::from_secs(2)), 221 ci_logs: None, 222 })); 223 224 let listener = TcpListener::bind(("127.0.0.1", port)).await.unwrap(); 225 let bound = listener.local_addr().unwrap().port(); 226 tokio::spawn(govern_decay()); 227 tokio::spawn(async move { 228 let _ = knot_ssh::serve_on_socket(listener, host_key, state).await; 229 }); 230 231 println!("READY pid={} port={}", std::process::id(), bound); 232 let mut reporter = tokio::time::interval(Duration::from_secs(1)); 233 tokio::select! { 234 _ = tokio::signal::ctrl_c() => {} 235 () = async { 236 loop { 237 reporter.tick().await; 238 println!("VmHWM {} MiB", vm_hwm_mib()); 239 } 240 } => {} 241 } 242 drop(scan); 243}