This repository has no description
1use std::collections::{BTreeSet, HashMap, HashSet};
2use std::fs;
3use std::path::{Path, PathBuf};
4use std::sync::{Arc, Mutex};
5use std::time::Duration;
6
7use axum::{Json, Router, routing::get};
8use base64::Engine;
9use base64::engine::general_purpose::STANDARD;
10use bytes::Bytes;
11use serde_json::json;
12use tempfile::TempDir;
13use url::Url;
14
15use knot_atproto::{Atproto, knot_did_document};
16use knot_cob::{CobHome, CobStore};
17use knot_cobs::{Grant, Registration, RegistryChange};
18use knot_events::{EventLog, GlobalSubscriberLimit, PerPeerSubscriberLimit, SubscriberGate};
19use knot_git::{
20 EntryKind, Identity, Layout, NewCommit, RefUpdate, Repo, StagedAction, StagedChange,
21};
22use knot_index::{Index, Resolved};
23use knot_runtime::{
24 Clock, Entropy, FakeHttp, HttpRequest, HttpResponse, K256Signer, ManualClock, NetworkError,
25 PublicKeyBytes, SeededEntropy, Signer, UnixMicros,
26};
27use knot_secrets::{MasterKey, SealedStore};
28use knot_types::{
29 AccountDid, AdmissionPolicy, AuthorName, BranchName, Email, KnotHostname, KnotId,
30 KnotServiceUrl, Oid, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, UnixSeconds,
31};
32use knot_xrpc::{
33 BlobReadBudget, BodyLimit, Budgets, ByteLimits, CobLocks, Committer, GlobalInflight,
34 GlobalQuota, LanguagesPushBudget, LanguagesReadBudget, LimitConfig, MaxWireBytes,
35 PerActorQuota, PerPeerInflight, PreAuthLimiter, ReadBudget, ReservationTtl, Reservations,
36 ResponseLimit, TreeReadBudget, XrpcState,
37};
38
39use crate::realdata::{
40 RealActors, SAMPLE_TAG, SAMPLED_REPOS, STRANGER_POOL, SUBJECT_POOL, admin_signer, owner_signer,
41};
42use crate::trace::{RepoCollaborators, RoundNumber, Snapshot};
43use crate::workload::Rng;
44
45const EMPTY_TREE: &str = "4b825dc642cb6eb9a060e54bf8d69288fbee4904";
46const KNOT_HOST: &str = "knot.nel.pet";
47const ADMIN_HOST: &str = "admin.nel.pet";
48const STRANGER_SEED_BASE: u64 = 1_000;
49const PDS_ENDPOINT: &str = "https://pds.nel.pet";
50const START_MICROS: u64 = 1_000_000_000;
51const NO_REFLOG_EXPIRY_FLOOR_SECS: i64 = i64::MAX / 4;
52const NO_PRUNE_EXPIRY_GRACE_SECS: u64 = (i64::MAX / 4) as u64;
53const MAINTAIN_ATTEMPTS: usize = 4;
54
55pub(crate) const SUBJECT_DIDS: [&str; 8] = [
56 "did:plc:limpet",
57 "did:plc:whelk",
58 "did:plc:mussel",
59 "did:plc:conch",
60 "did:plc:scallop",
61 "did:plc:cuttle",
62 "did:plc:periwinkle",
63 "did:plc:nautilus",
64];
65
66pub(crate) type Responder =
67 Box<dyn Fn(&HttpRequest) -> Result<HttpResponse, NetworkError> + Send + Sync>;
68
69pub(crate) struct Actor {
70 pub host: KnotHostname,
71 pub did: AccountDid,
72 pub signer: K256Signer,
73}
74
75#[derive(Default)]
76pub(crate) struct Faults {
77 dropped: Mutex<HashSet<KnotHostname>>,
78}
79
80impl Faults {
81 fn is_dropped(&self, host: &KnotHostname) -> bool {
82 self.dropped.lock().expect("faults lock").contains(host)
83 }
84
85 pub(crate) fn drop_host(&self, host: &KnotHostname) {
86 self.dropped
87 .lock()
88 .expect("faults lock")
89 .insert(host.clone());
90 }
91
92 pub(crate) fn clear_host(&self, host: &KnotHostname) {
93 self.dropped.lock().expect("faults lock").remove(host);
94 }
95}
96
97pub(crate) struct Harness {
98 _dir: Option<TempDir>,
99 pub clock: Arc<ManualClock>,
100 pub faults: Arc<Faults>,
101 pub layout: Layout,
102 pub index: Arc<Index>,
103 pub knot_aud: KnotId,
104 router: Router,
105 pub admin: Actor,
106 pub strangers: Vec<Actor>,
107 pub subjects: Vec<AccountDid>,
108 pub seed_repo: RepoDid,
109 options: knot_maintenance::Options,
110 maintain_lock: Mutex<()>,
111}
112
113impl Harness {
114 pub(crate) fn build(seed: u64, stranger_pool: usize) -> Self {
115 let dir = tempfile::tempdir().expect("sim tempdir");
116 let knot = KnotId::new(format!("did:web:{KNOT_HOST}")).expect("knot did");
117 let layout = Layout::new(dir.path().join("repos"))
118 .reserving_meta(&knot)
119 .expect("reserve meta");
120 layout.bootstrap_meta(&knot).expect("bootstrap meta");
121 let meta_path = layout.meta_path(&knot).expect("meta path");
122
123 let entropy = Arc::new(SeededEntropy::new(seed ^ 0x5eed_0050));
124 let secrets = Arc::new(
125 SealedStore::open(
126 dir.path().join("keys.sealed"),
127 &MasterKey::new([7u8; 32]).unwrap(),
128 Box::new(SeededEntropy::new(seed ^ 0x5ec0)),
129 )
130 .expect("sealed store"),
131 );
132 let knot_pubkey = secrets.ensure(&knot).expect("seal knot key");
133 let knot_service_url =
134 KnotServiceUrl::new(format!("https://{KNOT_HOST}")).expect("knot service url");
135
136 let admin = actor(ADMIN_HOST, 1);
137 let strangers: Vec<Actor> = (0..stranger_pool)
138 .map(|index| {
139 let host = format!("stranger{index}.nel.pet");
140 actor(&host, STRANGER_SEED_BASE + index as u64)
141 })
142 .collect();
143 let subjects: Vec<AccountDid> = SUBJECT_DIDS
144 .iter()
145 .map(|did| AccountDid::new(*did).expect("subject did"))
146 .collect();
147
148 let seed_repo = RepoDid::new("did:plc:squid").expect("seed repo did");
149 seed_on_disk(&layout, &seed_repo);
150 register_seed_repo(&meta_path, &knot, &secrets, &admin.did, &seed_repo);
151
152 let index = Arc::new(Index::new(meta_path.clone(), layout.clone()));
153 index.rebuild().expect("rebuild index");
154 index.warm_collaborators();
155
156 let clock = Arc::new(ManualClock::new(UnixMicros::new(START_MICROS)));
157 let pubkeys = pubkey_map(&admin, &strangers);
158 let faults = Arc::new(Faults::default());
159 let plc = Url::parse("https://plc.directory/").expect("plc url");
160 let responder = build_responder(Arc::clone(&faults), pubkeys, HashMap::new(), plc.clone());
161 let knot_aud = knot.clone();
162 let did_document = knot_did_document(&knot, &knot_pubkey, &knot_service_url);
163 let router = assemble_router(StateParts {
164 layout: layout.clone(),
165 index: Arc::clone(&index),
166 responder,
167 secrets,
168 entropy: entropy as Arc<dyn Entropy>,
169 admins: BTreeSet::from([admin.did.clone()]),
170 admission: AdmissionPolicy::Closed,
171 knot: knot.clone(),
172 knot_hostname: KnotHostname::new(KNOT_HOST).unwrap(),
173 meta_path,
174 knot_service_url,
175 service_owner: admin.did.clone(),
176 clock: Arc::clone(&clock),
177 plc,
178 did_document,
179 });
180 let options = maintenance_options();
181
182 Self {
183 _dir: Some(dir),
184 clock,
185 faults,
186 layout,
187 index,
188 knot_aud,
189 router,
190 admin,
191 strangers,
192 subjects,
193 seed_repo,
194 options,
195 maintain_lock: Mutex::new(()),
196 }
197 }
198
199 pub fn open(target: &Path, seed: u64) -> Result<(Self, RealActors), OpenError> {
200 let config =
201 knot_config::load(Some(&target.join("config.toml"))).map_err(OpenError::Config)?;
202 let hostname =
203 KnotHostname::new(config.server.hostname.clone()).map_err(|_| OpenError::Hostname)?;
204 let knot = hostname.knot_did();
205 let object_format = config.object_format().ok_or(OpenError::ObjectFormat)?;
206 let default_branch = BranchName::new(config.repo.default_branch.as_str())
207 .map_err(|_| OpenError::DefaultBranch)?;
208 let admin_did = config
209 .server
210 .admins
211 .first()
212 .cloned()
213 .ok_or(OpenError::NoAdmin)?;
214 let admission = config.acl.admission;
215 let plc = config.atproto.plc_directory.clone();
216
217 let master_key_env = config.secrets.master_key_env.clone();
218 let master_key = MasterKey::new(
219 STANDARD
220 .decode(
221 std::env::var(&master_key_env)
222 .map_err(|_| OpenError::MasterKeyEnv(master_key_env.clone()))?
223 .trim(),
224 )
225 .map_err(|_| OpenError::MasterKeyDecode)?,
226 )?;
227
228 let scratch = materialize_scratch(target, &knot)?;
229 let root = scratch.path().to_path_buf();
230
231 let layout = Layout::new(root.join("repos"))
232 .with_default_branch(default_branch)
233 .with_object_format(object_format)
234 .reserving_meta(&knot)?;
235 let meta_path = layout.meta_path(&knot)?;
236 let secrets = Arc::new(SealedStore::open(
237 root.join("sealed-keys"),
238 &master_key,
239 Box::new(SeededEntropy::new(seed ^ 0x5ec0)),
240 )?);
241 let knot_pubkey = secrets.ensure(&knot)?;
242 let knot_service_url =
243 KnotServiceUrl::new(format!("https://{hostname}")).map_err(|_| OpenError::Hostname)?;
244
245 let index = Arc::new(Index::new(meta_path.clone(), layout.clone()));
246 index.rebuild()?;
247 index.warm_collaborators();
248
249 let mut hosted = index.hosted_repos();
250 hosted.sort();
251 let rng = Rng::new(seed ^ SAMPLE_TAG);
252 let owner_repos: Vec<(RepoDid, OwnerDid)> = SAMPLED_REPOS
253 .sample(&hosted, &rng)
254 .into_iter()
255 .filter_map(|repo| match index.owner_of(&repo) {
256 Resolved::Ready(Some(owner)) => Some((repo, owner)),
257 _ => None,
258 })
259 .collect();
260 let subjects = SUBJECT_POOL.sample(&distinct_members(&index), &rng);
261
262 let admin = Actor {
263 host: hostname.clone(),
264 did: admin_did.clone(),
265 signer: admin_signer(seed),
266 };
267 let strangers: Vec<Actor> = (0..STRANGER_POOL.count())
268 .map(|index| {
269 actor(
270 &format!("stranger{index}.nel.pet"),
271 STRANGER_SEED_BASE + index as u64,
272 )
273 })
274 .collect();
275
276 let mut did_overrides: HashMap<AccountDid, PublicKeyBytes> = HashMap::new();
277 did_overrides.insert(admin_did.clone(), admin.signer.public_key());
278 owner_repos.iter().for_each(|(_, owner)| {
279 did_overrides
280 .entry(AccountDid::from(owner.clone()))
281 .or_insert_with(|| owner_signer(seed, owner).public_key());
282 });
283
284 let seed_repo = owner_repos
285 .first()
286 .map(|(repo, _)| repo.clone())
287 .or_else(|| hosted.first().cloned())
288 .ok_or(OpenError::NoRepos)?;
289
290 let clock = Arc::new(ManualClock::new(UnixMicros::new(START_MICROS)));
291 let faults = Arc::new(Faults::default());
292 let responder = build_responder(
293 Arc::clone(&faults),
294 pubkey_map(&admin, &strangers),
295 did_overrides,
296 plc.clone(),
297 );
298 let entropy = Arc::new(SeededEntropy::new(seed ^ 0x5eed_0050));
299 let knot_aud = knot.clone();
300 let did_document = knot_did_document(&knot, &knot_pubkey, &knot_service_url);
301 let router = assemble_router(StateParts {
302 layout: layout.clone(),
303 index: Arc::clone(&index),
304 responder,
305 secrets,
306 entropy: entropy as Arc<dyn Entropy>,
307 admins: BTreeSet::from([admin_did.clone()]),
308 admission,
309 knot: knot.clone(),
310 knot_hostname: hostname,
311 meta_path,
312 knot_service_url,
313 service_owner: admin_did.clone(),
314 clock: Arc::clone(&clock),
315 plc,
316 did_document,
317 });
318
319 let harness = Self {
320 _dir: Some(scratch),
321 clock,
322 faults,
323 layout,
324 index,
325 knot_aud,
326 router,
327 admin,
328 strangers,
329 subjects,
330 seed_repo,
331 options: maintenance_options(),
332 maintain_lock: Mutex::new(()),
333 };
334 Ok((harness, RealActors { owner_repos, seed }))
335 }
336
337 pub(crate) fn router(&self) -> Router {
338 self.router.clone()
339 }
340
341 pub(crate) fn now_seconds(&self) -> UnixSeconds {
342 UnixSeconds::new((self.clock.now_unix_micros().get() / 1_000_000) as i64)
343 }
344
345 pub(crate) fn advance(&self, delta: std::time::Duration) {
346 self.clock.advance(delta);
347 }
348
349 pub(crate) fn maintain(&self, repo: &RepoDid) -> Result<(), String> {
350 let _serialized = self.maintain_lock.lock().expect("maintenance lock");
351 let now = self.now_seconds();
352 let attempt = || -> Result<(), knot_maintenance::MaintError> {
353 let opened = self
354 .layout
355 .open(repo)
356 .map_err(knot_maintenance::MaintError::from)?;
357 knot_maintenance::run_repo(&opened, now, &self.options).map(|_| ())
358 };
359 (1..MAINTAIN_ATTEMPTS)
360 .fold(attempt(), |result, _| result.or_else(|_| attempt()))
361 .map_err(|error| error.to_string())
362 }
363
364 pub(crate) fn populate(&self, repo: &RepoDid) {
365 let opened = self.layout.open(repo).expect("open created repo");
366 write_history(&opened, repo.as_str());
367 }
368
369 pub(crate) fn snapshot(&self, round: RoundNumber, repos: &[RepoDid]) -> Snapshot {
370 let collaborators = repos
371 .iter()
372 .map(|repo| {
373 let _ = self.index.ensure_collaborators(repo);
374 RepoCollaborators {
375 repo: repo.clone(),
376 subjects: sorted(grant_subjects(self.index.collaborator_entries(repo))),
377 }
378 })
379 .collect();
380 let repo_list = {
381 let mut repos: Vec<RepoDid> = self.index.hosted_repos().to_vec();
382 repos.sort();
383 repos
384 };
385 Snapshot {
386 round,
387 clock_micros: self.clock.now_unix_micros(),
388 members: sorted(grant_subjects(self.index.member_entries())),
389 blocked: sorted(grant_subjects(self.index.blocked_entries())),
390 repos: repo_list,
391 collaborators,
392 }
393 }
394}
395
396fn actor(host: &str, seed: u64) -> Actor {
397 Actor {
398 host: KnotHostname::new(host).expect("actor hostname"),
399 did: AccountDid::new(format!("did:web:{host}")).expect("actor did"),
400 signer: K256Signer::generate(&SeededEntropy::new(seed)),
401 }
402}
403
404fn pubkey_map(admin: &Actor, strangers: &[Actor]) -> HashMap<KnotHostname, PublicKeyBytes> {
405 std::iter::once((admin.host.clone(), admin.signer.public_key()))
406 .chain(
407 strangers
408 .iter()
409 .map(|actor| (actor.host.clone(), actor.signer.public_key())),
410 )
411 .collect()
412}
413
414fn build_responder(
415 faults: Arc<Faults>,
416 pubkeys: HashMap<KnotHostname, PublicKeyBytes>,
417 did_overrides: HashMap<AccountDid, PublicKeyBytes>,
418 plc: Url,
419) -> Responder {
420 let plc_signer = K256Signer::generate(&SeededEntropy::new(7));
421 Box::new(move |request: &HttpRequest| {
422 if request.method == http::Method::POST {
423 return Ok(ok_body(Bytes::new()));
424 }
425 let not_found = || HttpResponse {
426 status: http::StatusCode::NOT_FOUND,
427 headers: http::HeaderMap::new(),
428 body: Bytes::new(),
429 };
430 let Ok(host) = KnotHostname::new(request.url.host_str().unwrap_or_default()) else {
431 return Ok(not_found());
432 };
433 if faults.is_dropped(&host) {
434 return Err(NetworkError::Timeout(
435 "identity resolution dropped by simulation".to_string(),
436 ));
437 }
438 let is_plc = request.url.host() == plc.host();
439 let requested_did = if is_plc {
440 request
441 .url
442 .path_segments()
443 .and_then(|mut segments| segments.rfind(|segment| !segment.is_empty()))
444 .and_then(|segment| AccountDid::new(segment).ok())
445 } else {
446 AccountDid::new(format!("did:web:{host}")).ok()
447 };
448 let Some(requested_did) = requested_did else {
449 return Ok(not_found());
450 };
451 if let Some(sec1) = did_overrides.get(&requested_did) {
452 return Ok(ok_body(did_doc(&requested_did, sec1)));
453 }
454 if is_plc {
455 return Ok(ok_body(did_doc(&requested_did, &plc_signer.public_key())));
456 }
457 match pubkeys.get(&host) {
458 Some(sec1) => Ok(ok_body(did_doc(&requested_did, sec1))),
459 None => Ok(not_found()),
460 }
461 })
462}
463
464fn did_doc(did: &AccountDid, sec1: &PublicKeyBytes) -> Bytes {
465 let did = did.as_str();
466 let multikey = knot_types::crypto::multikey(0xe7, sec1.as_bytes());
467 Bytes::from(
468 serde_json::to_vec(&json!({
469 "id": did,
470 "alsoKnownAs": [],
471 "verificationMethod": [{
472 "id": format!("{did}#atproto"),
473 "type": "Multikey",
474 "controller": did,
475 "publicKeyMultibase": multikey
476 }],
477 "service": [{
478 "id": "#atproto_pds",
479 "type": "AtprotoPersonalDataServer",
480 "serviceEndpoint": PDS_ENDPOINT
481 }]
482 }))
483 .expect("did doc serializes"),
484 )
485}
486
487fn ok_body(body: Bytes) -> HttpResponse {
488 HttpResponse {
489 status: http::StatusCode::OK,
490 headers: http::HeaderMap::new(),
491 body,
492 }
493}
494
495fn seed_on_disk(layout: &Layout, did: &RepoDid) {
496 let repo = layout.create(did).expect("create seed repo");
497 write_history(&repo, "reef");
498}
499
500fn write_history(repo: &Repo, marker: &str) {
501 let identity = Identity {
502 name: AuthorName::new("nel"),
503 email: Email::new("nel@oyster.cafe"),
504 time: UnixSeconds::new(1_700_000_000),
505 offset_seconds: 0,
506 };
507 let main = RefName::new("refs/heads/main").expect("main ref");
508 let first_tree = repo
509 .write_staged_tree(
510 Oid::from_hex(EMPTY_TREE).expect("empty tree"),
511 &[StagedChange {
512 path: knot_types::RepoPath::new("README.md").unwrap(),
513 action: StagedAction::Put {
514 content: format!("# {marker}\n").into_bytes(),
515 kind: EntryKind::Blob,
516 },
517 }],
518 )
519 .expect("first tree");
520 let root = repo
521 .write_commit(&NewCommit {
522 tree: first_tree,
523 parents: Vec::new(),
524 author: identity.clone(),
525 committer: identity.clone(),
526 message: "root".to_string(),
527 extra_headers: Vec::new(),
528 })
529 .expect("root commit");
530 let second_tree = repo
531 .write_staged_tree(
532 first_tree,
533 &[StagedChange {
534 path: knot_types::RepoPath::new("src/main.rs").unwrap(),
535 action: StagedAction::Put {
536 content: b"fn main() {}\n".to_vec(),
537 kind: EntryKind::Blob,
538 },
539 }],
540 )
541 .expect("second tree");
542 let tip = repo
543 .write_commit(&NewCommit {
544 tree: second_tree,
545 parents: vec![root],
546 author: identity.clone(),
547 committer: identity,
548 message: "add main".to_string(),
549 extra_headers: Vec::new(),
550 })
551 .expect("second commit");
552 repo.update_ref(&RefUpdate::Create {
553 name: main.clone(),
554 new: tip,
555 })
556 .expect("create main");
557 repo.set_head(&main).expect("set head");
558}
559
560fn register_seed_repo(
561 meta_path: &std::path::Path,
562 knot: &KnotId,
563 secrets: &SealedStore,
564 owner: &AccountDid,
565 repo: &RepoDid,
566) {
567 let meta = Repo::open(meta_path).expect("open meta");
568 let store = CobStore::new(&meta);
569 let signer = secrets.signer(knot).expect("knot signer");
570 store
571 .create(
572 &CobHome::from(knot),
573 &RegistryChange::Register(Registration {
574 owner: OwnerDid::new(owner.as_str()).expect("owner did"),
575 rkey: RepoRkey::new("anemone").expect("rkey"),
576 name: RepoName::new("anemone").expect("name"),
577 repo: repo.clone(),
578 created_at: UnixSeconds::new(1),
579 }),
580 &signer,
581 UnixSeconds::new(1),
582 )
583 .expect("register seed repo");
584}
585
586fn grant_subjects(resolved: Resolved<Vec<Grant>>) -> Vec<AccountDid> {
587 match resolved {
588 Resolved::Ready(grants) => grants
589 .into_iter()
590 .map(|grant| grant.subject.clone())
591 .collect(),
592 Resolved::Warming => Vec::new(),
593 }
594}
595
596fn sorted<T: Ord>(mut values: Vec<T>) -> Vec<T> {
597 values.sort();
598 values
599}
600
601fn distinct_members(index: &Index) -> Vec<AccountDid> {
602 let mut members = grant_subjects(index.member_entries());
603 members.sort();
604 members.dedup();
605 members
606}
607
608fn materialize_scratch(target: &Path, knot: &KnotId) -> Result<TempDir, OpenError> {
609 let base = target
610 .parent()
611 .filter(|parent| !parent.as_os_str().is_empty())
612 .map(Path::to_path_buf)
613 .unwrap_or_else(|| PathBuf::from("."));
614 let source_repos = target.join("repos");
615 let source_meta = Layout::new(source_repos.clone()).meta_path(knot)?;
616 let scratch = tempfile::tempdir_in(&base).map_err(OpenError::Scratch)?;
617 let root = scratch.path();
618 let scratch_repos = root.join("repos");
619 let scratch_meta = Layout::new(scratch_repos.clone()).meta_path(knot)?;
620 fs::copy(target.join("sealed-keys"), root.join("sealed-keys")).map_err(OpenError::Scratch)?;
621 hardlink_tree(&source_repos, &scratch_repos).map_err(OpenError::Scratch)?;
622 fs::remove_dir_all(&scratch_meta).map_err(OpenError::Scratch)?;
623 copy_tree(&source_meta, &scratch_meta).map_err(OpenError::Scratch)?;
624 Ok(scratch)
625}
626
627fn hardlink_tree(src: &Path, dst: &Path) -> std::io::Result<()> {
628 fs::create_dir_all(dst)?;
629 fs::read_dir(src)?.try_for_each(|entry| {
630 let entry = entry?;
631 let from = entry.path();
632 let to = dst.join(entry.file_name());
633 if entry.file_type()?.is_dir() {
634 hardlink_tree(&from, &to)
635 } else {
636 fs::hard_link(&from, &to).map(|_| ())
637 }
638 })
639}
640
641fn copy_tree(src: &Path, dst: &Path) -> std::io::Result<()> {
642 fs::create_dir_all(dst)?;
643 fs::read_dir(src)?.try_for_each(|entry| {
644 let entry = entry?;
645 let from = entry.path();
646 let to = dst.join(entry.file_name());
647 if entry.file_type()?.is_dir() {
648 copy_tree(&from, &to)
649 } else {
650 fs::copy(&from, &to).map(|_| ())
651 }
652 })
653}
654
655fn maintenance_options() -> knot_maintenance::Options {
656 knot_maintenance::Options {
657 repack_max_objects: knot_maintenance::ObjectCount::new(1_000_000),
658 geometric_factor: knot_maintenance::GeometricFactor::full_repack(),
659 prune_grace: knot_maintenance::PruneGrace::from_secs(NO_PRUNE_EXPIRY_GRACE_SECS),
660 reflog_floor: knot_maintenance::ReflogRetention::from_secs(
661 NO_REFLOG_EXPIRY_FLOOR_SECS as u64,
662 ),
663 commit_graph: true,
664 multi_pack_index: true,
665 bitmap: true,
666 }
667}
668
669struct StateParts {
670 layout: Layout,
671 index: Arc<Index>,
672 responder: Responder,
673 secrets: Arc<SealedStore>,
674 entropy: Arc<dyn Entropy>,
675 admins: BTreeSet<AccountDid>,
676 admission: AdmissionPolicy,
677 knot: KnotId,
678 knot_hostname: KnotHostname,
679 meta_path: PathBuf,
680 knot_service_url: KnotServiceUrl,
681 service_owner: AccountDid,
682 clock: Arc<ManualClock>,
683 plc: Url,
684 did_document: serde_json::Value,
685}
686
687fn assemble_router(parts: StateParts) -> Router {
688 let StateParts {
689 layout,
690 index,
691 responder,
692 secrets,
693 entropy,
694 admins,
695 admission,
696 knot,
697 knot_hostname,
698 meta_path,
699 knot_service_url,
700 service_owner,
701 clock,
702 plc,
703 did_document,
704 } = parts;
705 let atproto = Arc::new(Atproto::new(
706 FakeHttp::new(responder),
707 Arc::clone(&clock),
708 knot.clone(),
709 knot_atproto::PlcDirectory::new(plc).expect("plc directory"),
710 ));
711 let state = Arc::new(XrpcState {
712 layout: layout.clone(),
713 index: Arc::clone(&index),
714 atproto,
715 secrets,
716 entropy,
717 ci_logs: None,
718 admins,
719 admission,
720 knot_did: knot,
721 knot_hostname,
722 meta_path,
723 knot_service_url,
724 limiter: Arc::new(PreAuthLimiter::with_config(LimitConfig {
725 rate: None,
726 per_peer_inflight: Some(PerPeerInflight::new(4096)),
727 global_inflight: Some(GlobalInflight::new(4096)),
728 })),
729 cob_locks: Arc::new(CobLocks::default()),
730 reservations: Arc::new(Reservations::new(
731 ReservationTtl::new(1_000_000),
732 PerActorQuota::new(256),
733 GlobalQuota::new(256),
734 )),
735 trusted_proxy_header: None,
736 committer: Committer {
737 name: AuthorName::new("knot"),
738 email: Email::new("knot@nel.pet"),
739 },
740 byte_limits: ByteLimits {
741 body: BodyLimit::new(256 * 1024),
742 response: ResponseLimit::new(8 * 1024 * 1024),
743 pack: MaxWireBytes::new(1024 * 1024 * 1024),
744 ..ByteLimits::default()
745 },
746 budgets: Budgets {
747 tree_last_commit: TreeReadBudget::new(ReadBudget::Unbounded),
748 blob_last_commit: BlobReadBudget::new(ReadBudget::Unbounded),
749 languages: LanguagesReadBudget::new(ReadBudget::Unbounded),
750 languages_push: LanguagesPushBudget::new(Duration::from_secs(120)),
751 },
752 git_http: Arc::new(FakeHttp::new(|_request: &HttpRequest| {
753 Err(NetworkError::Connect(
754 "simulation serves no git upstream".to_string(),
755 ))
756 })),
757 pack_limits: knot_pack::PackLimits::default(),
758 service_owner,
759 events: Arc::new(EventLog::new(
760 Arc::clone(&clock),
761 knot_events::ReplayBounds::new(
762 knot_events::ReplayEvents::new(4096).expect("replay event maximum is nonzero"),
763 knot_events::ReplayBytes::new(64 << 20).expect("replay byte maximum is nonzero"),
764 ),
765 )),
766 subscriber_gate: Arc::new(SubscriberGate::new(
767 GlobalSubscriberLimit::new(256),
768 PerPeerSubscriberLimit::new(64),
769 )),
770 maintenance: knot_maintenance::MaintenanceHandle::disabled(),
771 appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(),
772 slots: knot_resource::Slots::testing(8),
773 lfs: None,
774 catalog: Arc::new(knot_messages::Catalog::defaults()),
775 });
776 let resolver: Arc<dyn knot_pack::RepoResolver> = {
777 let index = Arc::clone(&index);
778 Arc::new(move |target: &knot_pack::RepoTarget| match target {
779 knot_pack::RepoTarget::Did(did) => match index.owner_of(did) {
780 Resolved::Ready(Some(_)) => knot_pack::RepoLookup::Hosted(did.clone()),
781 Resolved::Ready(None) => knot_pack::RepoLookup::Unhosted,
782 Resolved::Warming => knot_pack::RepoLookup::Unavailable,
783 },
784 knot_pack::RepoTarget::OwnerPath(owner, path) => {
785 match index.resolve_clone_path(owner, path) {
786 Resolved::Ready(Some(found)) => knot_pack::RepoLookup::Hosted(found),
787 Resolved::Ready(None) => knot_pack::RepoLookup::Unhosted,
788 Resolved::Warming => knot_pack::RepoLookup::Unavailable,
789 }
790 }
791 })
792 };
793 knot_pack::router(layout, resolver, Arc::clone(&clock) as Arc<dyn Clock>)
794 .merge(knot_xrpc::router(Arc::clone(&state)))
795 .route(
796 "/.well-known/did.json",
797 get(move || {
798 let document = did_document.clone();
799 async move { Json(document) }
800 }),
801 )
802}
803
804#[derive(Debug)]
805pub enum OpenError {
806 MasterKeyEnv(String),
807 MasterKeyDecode,
808 NoAdmin,
809 NoRepos,
810 Scratch(std::io::Error),
811 Hostname,
812 ObjectFormat,
813 DefaultBranch,
814 Config(knot_config::LoadError),
815 Git(knot_git::GitError),
816 Secrets(knot_secrets::SecretsError),
817 Index(knot_index::IndexError),
818}
819
820impl std::fmt::Display for OpenError {
821 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
822 match self {
823 Self::MasterKeyEnv(name) => write!(f, "{name} is not set"),
824 Self::MasterKeyDecode => write!(f, "master key is not valid base64"),
825 Self::NoAdmin => write!(f, "config lists no admin"),
826 Self::NoRepos => write!(f, "target hosts no repos to sample"),
827 Self::Scratch(error) => {
828 write!(f, "materialize disposable scratch copy of target: {error}")
829 }
830 Self::Hostname => write!(f, "config hostname is not a valid knot hostname"),
831 Self::ObjectFormat => {
832 write!(f, "config git.object_format is not a valid object format")
833 }
834 Self::DefaultBranch => write!(f, "config default branch is not a valid branch name"),
835 Self::Config(error) => write!(f, "read config: {error}"),
836 Self::Git(error) => write!(f, "open target repos: {error}"),
837 Self::Secrets(error) => write!(f, "open sealed key store: {error}"),
838 Self::Index(error) => write!(f, "rebuild index: {error}"),
839 }
840 }
841}
842
843impl std::error::Error for OpenError {}
844
845impl From<knot_git::GitError> for OpenError {
846 fn from(error: knot_git::GitError) -> Self {
847 Self::Git(error)
848 }
849}
850
851impl From<knot_secrets::SecretsError> for OpenError {
852 fn from(error: knot_secrets::SecretsError) -> Self {
853 Self::Secrets(error)
854 }
855}
856
857impl From<knot_index::IndexError> for OpenError {
858 fn from(error: knot_index::IndexError) -> Self {
859 Self::Index(error)
860 }
861}