This repository has no description
0

Configure Feed

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

core / knot2 / crates / knot-sim / src / harness.rs
30 kB 861 lines
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}