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 / realdata.rs
16 kB 536 lines
1use std::path::Path; 2use std::sync::Arc; 3 4use axum::body::Bytes; 5use futures::future::join_all; 6use futures::stream::StreamExt; 7use http::Method; 8use serde_json::json; 9 10use knot_runtime::{K256Signer, SeededEntropy}; 11use knot_types::{AccountDid, HttpStatus, KnotHostname, OwnerDid, RepoDid}; 12 13use crate::harness::Harness; 14use crate::trace::{OperationIndex, Outcome, RoundNumber, Snapshot, Step, Trace, fnv1a}; 15use crate::workload::{ 16 Request, Rng, body_digest, drive_kill, enc, encode_body, http_call, jwt_window, mint, 17 repo_did_of, 18}; 19 20pub(crate) const SAMPLE_TAG: u64 = 0x5a3d_e005; 21pub(crate) const SAMPLED_REPOS: SampledRepos = SampledRepos(24); 22pub(crate) const SUBJECT_POOL: SubjectPool = SubjectPool(16); 23pub(crate) const STRANGER_POOL: StrangerPool = StrangerPool(32); 24const ADMIN_SIGNER_TAG: u64 = 0x0ad0_0001; 25const DRIVE_TAG: u64 = 0x11c7_0005; 26 27#[derive(Clone, Copy)] 28pub(crate) struct SampledRepos(usize); 29 30#[derive(Clone, Copy)] 31pub(crate) struct SubjectPool(usize); 32 33#[derive(Clone, Copy)] 34pub(crate) struct StrangerPool(usize); 35 36impl SampledRepos { 37 pub(crate) fn sample(self, source: &[RepoDid], rng: &Rng) -> Vec<RepoDid> { 38 pick(source, self.0, rng) 39 } 40} 41 42impl SubjectPool { 43 pub(crate) fn sample(self, source: &[AccountDid], rng: &Rng) -> Vec<AccountDid> { 44 pick(source, self.0, rng) 45 } 46} 47 48impl StrangerPool { 49 pub(crate) fn count(self) -> usize { 50 self.0 51 } 52} 53 54pub struct RealActors { 55 pub(crate) owner_repos: Vec<(RepoDid, OwnerDid)>, 56 pub(crate) seed: u64, 57} 58 59pub(crate) fn owner_signer(seed: u64, owner: &OwnerDid) -> K256Signer { 60 K256Signer::generate(&SeededEntropy::new(seed ^ fnv1a(owner.as_str().as_bytes()))) 61} 62 63pub(crate) fn admin_signer(seed: u64) -> K256Signer { 64 K256Signer::generate(&SeededEntropy::new(seed ^ ADMIN_SIGNER_TAG)) 65} 66 67fn pick<T: Clone>(source: &[T], count: usize, rng: &Rng) -> Vec<T> { 68 if source.is_empty() { 69 return Vec::new(); 70 } 71 (0..count.min(source.len())) 72 .scan(Vec::<usize>::new(), |taken, _| { 73 let start = rng.below(source.len() as u64) as usize; 74 let index = (start..source.len()) 75 .chain(0..start) 76 .find(|candidate| !taken.contains(candidate)) 77 .expect("distinct index exists while count <= len"); 78 taken.push(index); 79 Some(source[index].clone()) 80 }) 81 .collect() 82} 83 84pub async fn run( 85 target: &Path, 86 seed: u64, 87 rounds: u32, 88) -> Result<Trace, crate::harness::OpenError> { 89 let (harness, actors) = Harness::open(target, seed)?; 90 Ok(drive(Arc::new(harness), actors, rounds).await) 91} 92 93struct OpResult { 94 step: Step, 95 created: Option<RepoDid>, 96} 97 98enum Exec { 99 Http { 100 op: &'static str, 101 actor: String, 102 fault: &'static str, 103 request: Request, 104 killed: bool, 105 capture_create: bool, 106 }, 107 Maintain { 108 repo: RepoDid, 109 }, 110} 111 112async fn drive(harness: Arc<Harness>, actors: RealActors, rounds: u32) -> Trace { 113 let seed = actors.seed; 114 let rng = Rng::new(seed ^ DRIVE_TAG); 115 let harness_ref = &harness; 116 let actors_ref = &actors; 117 let rng_ref = &rng; 118 let base: Vec<RepoDid> = actors 119 .owner_repos 120 .iter() 121 .map(|(repo, _)| repo.clone()) 122 .collect(); 123 let (_created, steps, snapshots) = futures::stream::iter(0..rounds) 124 .fold( 125 ( 126 Vec::<RepoDid>::new(), 127 Vec::<Step>::new(), 128 Vec::<Snapshot>::new(), 129 ), 130 |(created, mut steps, mut snapshots), round_index| { 131 let harness = Arc::clone(harness_ref); 132 let base = &base; 133 async move { 134 let round = RoundNumber::new(round_index); 135 let working: Vec<RepoDid> = 136 base.iter().chain(created.iter()).cloned().collect(); 137 let execs = if round_index % 2 == 0 { 138 mutate(&harness, actors_ref, &working, rng_ref, round) 139 } else { 140 reads(&working, rng_ref) 141 }; 142 let dropped = arm_drops(&harness, &execs); 143 144 let tasks = execs 145 .into_iter() 146 .enumerate() 147 .map(|(position, exec)| { 148 let harness = Arc::clone(&harness); 149 let index = OperationIndex::new(position as u32); 150 tokio::spawn(async move { run_op(&harness, round, index, exec).await }) 151 }) 152 .collect::<Vec<_>>(); 153 let results: Vec<OpResult> = join_all(tasks) 154 .await 155 .into_iter() 156 .map(|joined| joined.expect("real-data op task mustn't panic")) 157 .collect(); 158 dropped 159 .iter() 160 .for_each(|host| harness.faults.clear_host(host)); 161 162 let fresh: Vec<RepoDid> = results 163 .iter() 164 .filter_map(|result| result.created.clone()) 165 .collect(); 166 fresh.iter().for_each(|did| harness.populate(did)); 167 steps.extend(results.into_iter().map(|result| result.step)); 168 169 let mut created = created; 170 created.extend(fresh); 171 harness.advance(round_advance(rng_ref)); 172 let snapshot_repos: Vec<RepoDid> = 173 base.iter().chain(created.iter()).cloned().collect(); 174 snapshots.push(harness.snapshot(round, &snapshot_repos)); 175 (created, steps, snapshots) 176 } 177 }, 178 ) 179 .await; 180 Trace { 181 seed, 182 steps, 183 snapshots, 184 } 185} 186 187fn round_advance(rng: &Rng) -> std::time::Duration { 188 std::time::Duration::from_micros(rng.below(3_000_000)) 189} 190 191fn arm_drops(harness: &Harness, execs: &[Exec]) -> Vec<KnotHostname> { 192 let hosts: Vec<KnotHostname> = execs 193 .iter() 194 .filter_map(|exec| match exec { 195 Exec::Http { 196 op: "probe", 197 fault: "drop_identity", 198 actor, 199 .. 200 } => KnotHostname::new(actor).ok(), 201 _ => None, 202 }) 203 .collect(); 204 hosts.iter().for_each(|host| harness.faults.drop_host(host)); 205 hosts 206} 207 208fn member_verb(draw: u64) -> (&'static str, &'static str) { 209 match draw { 210 0 => ("sh.tangled.knot.addMember", "addMember"), 211 1 => ("sh.tangled.knot.removeMember", "removeMember"), 212 2 => ("sh.tangled.knot.ban", "ban"), 213 _ => ("sh.tangled.knot.unban", "unban"), 214 } 215} 216 217struct Caller<'a> { 218 signer: &'a K256Signer, 219 did: &'a AccountDid, 220} 221 222fn signed_post( 223 harness: &Harness, 224 caller: Caller<'_>, 225 nsid: &'static str, 226 body: serde_json::Value, 227 skew: bool, 228 round: RoundNumber, 229 slot: usize, 230) -> Request { 231 let token = mint( 232 caller.signer, 233 caller.did, 234 &harness.knot_aud, 235 nsid, 236 jwt_window(harness, skew), 237 round, 238 OperationIndex::new(slot as u32), 239 ); 240 Request { 241 method: Method::POST, 242 uri: format!("/xrpc/{nsid}"), 243 token: Some(token), 244 body: encode_body(body), 245 actor: String::new(), 246 } 247} 248 249fn mutate( 250 harness: &Harness, 251 actors: &RealActors, 252 working: &[RepoDid], 253 rng: &Rng, 254 round: RoundNumber, 255) -> Vec<Exec> { 256 let subjects = &harness.subjects; 257 let mut execs: Vec<Exec> = Vec::new(); 258 259 subjects 260 .iter() 261 .filter(|_| rng.chance(2, 3)) 262 .for_each(|subject| { 263 let (nsid, short) = member_verb(rng.below(4)); 264 let skew = rng.chance(1, 5); 265 let request = signed_post( 266 harness, 267 Caller { 268 signer: &harness.admin.signer, 269 did: &harness.admin.did, 270 }, 271 nsid, 272 json!({ "subject": subject.as_str() }), 273 skew, 274 round, 275 execs.len(), 276 ); 277 execs.push(Exec::Http { 278 op: short, 279 actor: format!("admin:{short}"), 280 fault: skew_fault(skew), 281 request, 282 killed: false, 283 capture_create: false, 284 }); 285 }); 286 287 let collaborated: Vec<RepoDid> = if subjects.is_empty() { 288 Vec::new() 289 } else { 290 actors 291 .owner_repos 292 .iter() 293 .filter(|_| rng.chance(1, 2)) 294 .map(|(repo, owner)| { 295 let subject = &subjects[rng.below(subjects.len() as u64) as usize]; 296 let skew = rng.chance(1, 6); 297 let request = signed_post( 298 harness, 299 Caller { 300 signer: &owner_signer(actors.seed, owner), 301 did: &AccountDid::from(owner.clone()), 302 }, 303 "sh.tangled.repo.addCollaborator", 304 json!({ "repo": repo.as_str(), "subject": subject.as_str() }), 305 skew, 306 round, 307 execs.len(), 308 ); 309 execs.push(Exec::Http { 310 op: "addCollaborator", 311 actor: format!("owner:{}", owner.as_str()), 312 fault: skew_fault(skew), 313 request, 314 killed: false, 315 capture_create: false, 316 }); 317 repo.clone() 318 }) 319 .collect() 320 }; 321 322 working 323 .iter() 324 .filter(|repo| !collaborated.contains(*repo)) 325 .filter(|_| rng.chance(1, 3)) 326 .for_each(|repo| execs.push(Exec::Maintain { repo: repo.clone() })); 327 328 (0..rng.below(3)).for_each(|key| { 329 let name = format!("sim-repo-{}-{}", round.get(), key); 330 let skew = rng.chance(1, 8); 331 let request = signed_post( 332 harness, 333 Caller { 334 signer: &harness.admin.signer, 335 did: &harness.admin.did, 336 }, 337 "sh.tangled.repo.create", 338 json!({ "rkey": name, "name": name }), 339 skew, 340 round, 341 execs.len(), 342 ); 343 execs.push(Exec::Http { 344 op: "createRepo", 345 actor: "admin:create".to_string(), 346 fault: skew_fault(skew), 347 request, 348 killed: false, 349 capture_create: !skew, 350 }); 351 }); 352 353 if !subjects.is_empty() { 354 let probes = rng.below(3) as usize; 355 let slots: Vec<usize> = (0..harness.strangers.len()).collect(); 356 pick(&slots, probes, rng).into_iter().for_each(|slot| { 357 let stranger = &harness.strangers[slot]; 358 let drop = rng.chance(1, 2); 359 let request = signed_post( 360 harness, 361 Caller { 362 signer: &stranger.signer, 363 did: &stranger.did, 364 }, 365 "sh.tangled.knot.addMember", 366 json!({ "subject": subjects[0].as_str() }), 367 false, 368 round, 369 execs.len(), 370 ); 371 execs.push(Exec::Http { 372 op: "probe", 373 actor: stranger.host.to_string(), 374 fault: if drop { "drop_identity" } else { "none" }, 375 request, 376 killed: false, 377 capture_create: false, 378 }); 379 }); 380 } 381 382 execs 383} 384 385fn skew_fault(skew: bool) -> &'static str { 386 if skew { "clock_skew" } else { "none" } 387} 388 389fn reads(working: &[RepoDid], rng: &Rng) -> Vec<Exec> { 390 let mut execs: Vec<Exec> = [ 391 ("version", "/xrpc/sh.tangled.knot.version"), 392 ("owner", "/xrpc/sh.tangled.owner"), 393 ("didJson", "/.well-known/did.json"), 394 ] 395 .into_iter() 396 .map(|(op, uri)| anon_read(op, uri, rng.chance(1, 5))) 397 .collect(); 398 399 working 400 .iter() 401 .filter(|_| rng.chance(2, 3)) 402 .for_each(|repo| { 403 let killed = rng.chance(1, 5); 404 let exec = match rng.below(7) { 405 0 => repo_read("branches", "repo", repo, killed), 406 1 => repo_read("log", "repo", repo, killed), 407 2 => repo_read("describeRepo", "repoDid", repo, killed), 408 3 => anon_read( 409 "infoRefs", 410 &format!("/{}/info/refs?service=git-upload-pack", repo.as_str()), 411 killed, 412 ), 413 4 => repo_read("tree", "repo", repo, killed), 414 5 => blob_read(repo, killed), 415 _ => repo_read("languages", "repo", repo, killed), 416 }; 417 execs.push(exec); 418 }); 419 execs 420} 421 422fn anon_read(op: &'static str, uri: &str, killed: bool) -> Exec { 423 Exec::Http { 424 op, 425 actor: "anon".to_string(), 426 fault: if killed { "killed" } else { "none" }, 427 request: Request { 428 method: Method::GET, 429 uri: uri.to_string(), 430 token: None, 431 body: Bytes::new(), 432 actor: "anon".to_string(), 433 }, 434 killed, 435 capture_create: false, 436 } 437} 438 439fn repo_read(op: &'static str, param: &str, repo: &RepoDid, killed: bool) -> Exec { 440 anon_read( 441 op, 442 &format!("/xrpc/sh.tangled.repo.{op}?{param}={}", enc(repo.as_str())), 443 killed, 444 ) 445} 446 447fn blob_read(repo: &RepoDid, killed: bool) -> Exec { 448 anon_read( 449 "blob", 450 &format!( 451 "/xrpc/sh.tangled.repo.blob?repo={}&path=README.md", 452 enc(repo.as_str()) 453 ), 454 killed, 455 ) 456} 457 458async fn run_op( 459 harness: &Harness, 460 round: RoundNumber, 461 index: OperationIndex, 462 exec: Exec, 463) -> OpResult { 464 match exec { 465 Exec::Maintain { repo } => { 466 let outcome = match harness.maintain(&repo) { 467 Ok(()) => Outcome::Answered { 468 status: HttpStatus::new(200), 469 body: 0, 470 }, 471 Err(message) => Outcome::Answered { 472 status: HttpStatus::new(500), 473 body: fnv1a(message.as_bytes()), 474 }, 475 }; 476 OpResult { 477 step: Step { 478 round, 479 index, 480 op: "maintain", 481 actor: "knot".to_string(), 482 fault: "none", 483 outcome, 484 }, 485 created: None, 486 } 487 } 488 Exec::Http { 489 op, 490 actor, 491 fault, 492 request, 493 killed, 494 capture_create, 495 } => { 496 if killed { 497 drive_kill(harness.router(), request.method, &request.uri, request.body).await; 498 return OpResult { 499 step: Step { 500 round, 501 index, 502 op, 503 actor, 504 fault: "killed", 505 outcome: Outcome::Killed, 506 }, 507 created: None, 508 }; 509 } 510 let (status, body) = http_call( 511 harness.router(), 512 request.method, 513 &request.uri, 514 request.token.as_deref(), 515 request.body, 516 ) 517 .await; 518 let created = (capture_create && status == HttpStatus::new(200)) 519 .then(|| repo_did_of(&body).expect("200 createRepo response carries a repoDid")); 520 OpResult { 521 step: Step { 522 round, 523 index, 524 op, 525 actor, 526 fault, 527 outcome: Outcome::Answered { 528 status, 529 body: body_digest(&body), 530 }, 531 }, 532 created, 533 } 534 } 535 } 536}