This repository has no description
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}