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