| 1 | //! A crowd on one shared notebook: a host, and guests that join it through a relay, half of |
| 2 | //! them typing on the same page. Reports how long keys and carets take to reach the others, |
| 3 | //! what the relay passes on, and what the relay, the host and the guests cost. |
| 4 | //! |
| 5 | //! `cargo run --release -p notebook --features live --example live_crowd -- RELAY_BINARY |
| 6 | //! [PEERS=30] [TYPISTS=PEERS/2] [SECONDS=60]`, the relay's binary built with |
| 7 | //! `cargo build --release -p relay`. The relay, the host and the guests each run as a process |
| 8 | //! of their own, so `ps` tells their costs apart. |
| 9 | //! |
| 10 | //! For watching a crowd in the app: `fixture FOLDER PARAGRAPHS` writes the notebook the |
| 11 | //! host shares, and `guests URL CODE PEERS SECONDS` has that many join the app's share of it |
| 12 | //! and move their carets about its page. |
| 13 | |
| 14 | use notebook::{ |
| 15 | Replica, |
| 16 | live::{ |
| 17 | Caret, Guid, Hello, Presence, Spot, |
| 18 | share::{self, Guest, Host, Sharing}, |
| 19 | }, |
| 20 | session::{Background, Notebook, Section}, |
| 21 | }; |
| 22 | use onestore::{ |
| 23 | ExGuid, |
| 24 | op::{Edit, Op, PageOp}, |
| 25 | }; |
| 26 | use std::{ |
| 27 | collections::HashMap, |
| 28 | io::{BufRead, BufReader, Read, Write}, |
| 29 | net::TcpStream, |
| 30 | path::Path, |
| 31 | process::{Child, Command, Stdio}, |
| 32 | sync::{ |
| 33 | Arc, Mutex, |
| 34 | atomic::{AtomicBool, Ordering}, |
| 35 | }, |
| 36 | thread, |
| 37 | time::{Duration, Instant}, |
| 38 | }; |
| 39 | |
| 40 | /// Keys each typist types a second, and how often observers look. |
| 41 | const KEYS_PER_SECOND: f64 = 5.0; |
| 42 | const LOOK: Duration = Duration::from_millis(2); |
| 43 | /// Guests that time what they see. |
| 44 | const OBSERVERS: usize = 3; |
| 45 | |
| 46 | fn until(what: &str, limit: Duration, done: impl Fn() -> bool) { |
| 47 | let deadline = Instant::now() + limit; |
| 48 | while !done() { |
| 49 | assert!(Instant::now() < deadline, "{what}"); |
| 50 | thread::sleep(Duration::from_millis(5)); |
| 51 | } |
| 52 | } |
| 53 | |
| 54 | fn hello(name: &str) -> Hello { |
| 55 | Hello::new(name.into(), None).unwrap() |
| 56 | } |
| 57 | |
| 58 | fn main() { |
| 59 | let args: Vec<String> = std::env::args().skip(1).collect(); |
| 60 | match args.first().map(String::as_str) { |
| 61 | Some("host") => return host(&args[1], args[2].parse().unwrap()), |
| 62 | Some("fixture") => { |
| 63 | fixture(Path::new(&args[1]), args[2].parse().unwrap()); |
| 64 | return; |
| 65 | } |
| 66 | Some("guests") => { |
| 67 | let seconds = Duration::from_secs(args[4].parse().unwrap()); |
| 68 | return guests(&args[1], &args[2], args[3].parse().unwrap(), seconds); |
| 69 | } |
| 70 | _ => {} |
| 71 | } |
| 72 | let binary = args.first().expect("the relay's binary"); |
| 73 | let number = |at: usize, default: usize| args.get(at).map_or(default, |n| n.parse().unwrap()); |
| 74 | let peers = number(1, 30); |
| 75 | let typists = number(2, peers / 2); |
| 76 | let seconds = number(3, 60) as u64; |
| 77 | crowd(binary, peers, typists, Duration::from_secs(seconds)); |
| 78 | } |
| 79 | |
| 80 | /// Writes the notebook `Garden` in `directory`: a section whose page has `paragraphs` |
| 81 | /// paragraphs, `Typed:` then `Note 1:` and on. Its folder. |
| 82 | fn fixture(directory: &Path, paragraphs: usize) -> std::path::PathBuf { |
| 83 | let folder = directory.join("Garden"); |
| 84 | std::fs::create_dir_all(&folder).unwrap(); |
| 85 | let mut image = onestore::create_section("Garden.one", "Typed:", "Fixture").unwrap(); |
| 86 | for n in (1..paragraphs).rev() { |
| 87 | let arena = onestore::Arena::default(); |
| 88 | let mut section = onestore::Section::open(&arena, image).unwrap(); |
| 89 | let (space, ..) = section.pages().unwrap()[0]; |
| 90 | let page = section.page(space).unwrap(); |
| 91 | let (text, _) = texts(&page).swap_remove(0); |
| 92 | let right = onestore::page::text::new_id().unwrap(); |
| 93 | let ops = [ |
| 94 | PageOp::Split { |
| 95 | text, |
| 96 | at: 6, |
| 97 | paragraph: onestore::page::text::new_id().unwrap(), |
| 98 | right, |
| 99 | lists: Vec::new(), |
| 100 | }, |
| 101 | PageOp::Text { |
| 102 | text: right, |
| 103 | range: 0..0, |
| 104 | with: format!("Note {n}:"), |
| 105 | }, |
| 106 | ]; |
| 107 | let edit = Edit { |
| 108 | at: 134_000_000_000_000_000, |
| 109 | ops: ops.into_iter().map(|op| Op::Page { space, op }).collect(), |
| 110 | }; |
| 111 | section.apply("Fixture", &edit).unwrap(); |
| 112 | section.seal().unwrap(); |
| 113 | image = section.image(); |
| 114 | } |
| 115 | std::fs::write(folder.join("Garden.one"), image).unwrap(); |
| 116 | folder |
| 117 | } |
| 118 | |
| 119 | /// The host's process: shares a fixture notebook of a page with `paragraphs` paragraphs |
| 120 | /// through `url`, says its code, and reports what the section grew to once its input closes. |
| 121 | fn host(url: &str, paragraphs: usize) { |
| 122 | let directory = tempfile::tempdir().unwrap(); |
| 123 | let folder = fixture(directory.path(), paragraphs); |
| 124 | let file = folder.join("Garden.one"); |
| 125 | let storage = Notebook::open(&folder, directory.path().join("cache")) |
| 126 | .unwrap() |
| 127 | .into_storage(); |
| 128 | let host = Host::start( |
| 129 | storage, |
| 130 | hello("Host"), |
| 131 | Sharing::new("").unwrap(), |
| 132 | "Garden", |
| 133 | None, |
| 134 | Some(url), |
| 135 | || {}, |
| 136 | |_| Ok(()), |
| 137 | ) |
| 138 | .unwrap(); |
| 139 | until("no code", Duration::from_secs(30), || { |
| 140 | host.code().is_some_and(|code| share::code(&code).is_some()) |
| 141 | }); |
| 142 | println!("code {}", host.code().unwrap()); |
| 143 | let _ = std::io::stdin().read_to_end(&mut Vec::new()); |
| 144 | let image = std::fs::read(&file).unwrap(); |
| 145 | if let Some(output) = std::env::var_os("SNOWBOUND_LIVE_EVIDENCE") { |
| 146 | let output = Path::new(&output); |
| 147 | std::fs::create_dir_all(output).unwrap(); |
| 148 | std::fs::copy(&file, output.join("Garden.one")).unwrap(); |
| 149 | } |
| 150 | let size = image.len(); |
| 151 | let arena = onestore::Arena::default(); |
| 152 | let mut section = onestore::Section::open(&arena, image).unwrap(); |
| 153 | let pages = section.pages().unwrap().len(); |
| 154 | let conflicts: usize = section |
| 155 | .conflicts() |
| 156 | .unwrap() |
| 157 | .iter() |
| 158 | .map(|(_, c)| c.len()) |
| 159 | .sum(); |
| 160 | println!("section {size} bytes, {pages} pages, {conflicts} conflict pages"); |
| 161 | } |
| 162 | |
| 163 | /// A guest as the app holds one: its share, the notebook's background, the section open |
| 164 | /// and its file's identity. |
| 165 | struct Open { |
| 166 | guest: Arc<Guest>, |
| 167 | _background: Background, |
| 168 | section: Section, |
| 169 | file: [u8; 16], |
| 170 | } |
| 171 | |
| 172 | fn join(name: &str, code: &str, url: &str, cache: &Path) -> Open { |
| 173 | let welcome = share::join(hello(name), code, "", None, Some(url)).unwrap(); |
| 174 | let guest = Guest::start( |
| 175 | hello(name), |
| 176 | welcome.share, |
| 177 | welcome.secret, |
| 178 | None, |
| 179 | Some(url), |
| 180 | || {}, |
| 181 | ) |
| 182 | .unwrap(); |
| 183 | until("no host", Duration::from_secs(120), || { |
| 184 | guest.host().is_some() |
| 185 | }); |
| 186 | let mut notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap(); |
| 187 | let background = Background::hosted(Arc::clone(&guest), || {}).unwrap(); |
| 188 | background.watch(notebook.replicas()); |
| 189 | let replica = notebook.replica_path("Garden.one").unwrap(); |
| 190 | std::fs::create_dir_all(replica.parent().unwrap()).unwrap(); |
| 191 | let image = notebook.read_section("Garden.one").unwrap(); |
| 192 | let file = onestore::Header::parse(&image[..1024]).unwrap().file_id; |
| 193 | let replica = Replica::open_or_create(&replica, None, || Ok(image)).unwrap(); |
| 194 | let section = |
| 195 | Section::resume_hosted("Garden.one".into(), replica, Arc::clone(&guest), || {}).unwrap(); |
| 196 | background.hold("Garden.one", &section); |
| 197 | Open { |
| 198 | guest, |
| 199 | _background: background, |
| 200 | section, |
| 201 | file, |
| 202 | } |
| 203 | } |
| 204 | |
| 205 | /// Where in the open page's `paragraph`th text, `offset` units in, `open`'s guest is. |
| 206 | fn presence(open: &Open, paragraph: usize, offset: u32) -> Option<Presence> { |
| 207 | let (space, texts) = texts_of(&open.section)?; |
| 208 | let (text, _) = texts.get(paragraph)?; |
| 209 | let at = Spot { |
| 210 | text: Guid { |
| 211 | guid: text.guid, |
| 212 | n: text.n, |
| 213 | }, |
| 214 | offset, |
| 215 | }; |
| 216 | Some(Presence { |
| 217 | section: Some(open.file), |
| 218 | page: Some(Guid { |
| 219 | guid: space.guid, |
| 220 | n: space.n, |
| 221 | }), |
| 222 | caret: Some(Caret { |
| 223 | anchor: at, |
| 224 | focus: at, |
| 225 | }), |
| 226 | }) |
| 227 | } |
| 228 | |
| 229 | /// `peers` guests join the share `code` names through `url` and, for `seconds`, put their |
| 230 | /// carets in the page's paragraphs, each moving every few seconds. |
| 231 | fn guests(url: &str, code: &str, peers: usize, seconds: Duration) { |
| 232 | let directory = tempfile::tempdir().unwrap(); |
| 233 | let names = [ |
| 234 | "Ada", "Grace", "Alan", "Barbara", "Edsger", "Frances", "Donald", "Margaret", "Ken", |
| 235 | "Radia", "Dennis", "Hedy", "Linus", "Karen", "John", "Sophie", |
| 236 | ]; |
| 237 | let joining: Vec<_> = (0..peers) |
| 238 | .map(|n| { |
| 239 | let name = format!("{} {n}", names[n % names.len()]); |
| 240 | let (code, url, cache) = ( |
| 241 | code.to_owned(), |
| 242 | url.to_owned(), |
| 243 | directory.path().join(n.to_string()), |
| 244 | ); |
| 245 | thread::sleep(Duration::from_millis(100)); |
| 246 | thread::spawn(move || join(&name, &code, &url, &cache)) |
| 247 | }) |
| 248 | .collect(); |
| 249 | let guests: Vec<Open> = joining |
| 250 | .into_iter() |
| 251 | .map(|joined| joined.join().unwrap()) |
| 252 | .collect(); |
| 253 | let began = Instant::now(); |
| 254 | let mut turn = 0; |
| 255 | while began.elapsed() < seconds { |
| 256 | for (n, open) in guests.iter().enumerate() { |
| 257 | let paragraphs = texts_of(&open.section).map_or(1, |(_, texts)| texts.len()); |
| 258 | if (n + turn) % 4 == 0 || turn == 0 { |
| 259 | let at = (n * 7 + turn) % paragraphs; |
| 260 | if let Some(presence) = presence(open, at, (turn % 5) as u32) { |
| 261 | open.guest.set_presence(presence); |
| 262 | } |
| 263 | } |
| 264 | } |
| 265 | turn += 1; |
| 266 | thread::sleep(Duration::from_secs(1)); |
| 267 | } |
| 268 | } |
| 269 | |
| 270 | /// The page's object space and its paragraphs' texts: each one's object and what it reads. |
| 271 | fn texts_of(section: &Section) -> Option<(ExGuid, Vec<(ExGuid, String)>)> { |
| 272 | let (space, ..) = *section.pages().ok()?.first()?; |
| 273 | let page = section.page(space).ok()?; |
| 274 | Some((space, texts(&page))) |
| 275 | } |
| 276 | |
| 277 | fn texts(page: &onestore::page::Page) -> Vec<(ExGuid, String)> { |
| 278 | let outlines = page.objects.iter().filter_map(|object| match object { |
| 279 | onestore::page::PageObject::Outline(outline) => Some(outline), |
| 280 | _ => None, |
| 281 | }); |
| 282 | outlines |
| 283 | .flat_map(|outline| &outline.paragraphs) |
| 284 | .filter_map(|p| p.text().map(|t| (t.id, t.text.text().to_owned()))) |
| 285 | .collect() |
| 286 | } |
| 287 | |
| 288 | /// The key typist `typist` types `n`th, unique to both. |
| 289 | fn key(typist: usize, n: usize) -> char { |
| 290 | char::from_u32(0x4E00 + (typist * 600 + n) as u32).unwrap() |
| 291 | } |
| 292 | |
| 293 | fn typist_of(key: char) -> Option<(usize, usize)> { |
| 294 | let at = (key as u32).checked_sub(0x4E00)? as usize; |
| 295 | Some((at / 600, at % 600)) |
| 296 | } |
| 297 | |
| 298 | /// Seconds of CPU and kilobytes resident `ps` reports for `pid`. |
| 299 | fn cost(pid: u32) -> (f64, u64) { |
| 300 | let output = Command::new("ps") |
| 301 | .args(["-o", "time=,rss=", "-p", &pid.to_string()]) |
| 302 | .output() |
| 303 | .unwrap(); |
| 304 | let text = String::from_utf8_lossy(&output.stdout); |
| 305 | let mut words = text.split_whitespace(); |
| 306 | let time = words.next().unwrap_or("0:0"); |
| 307 | let seconds = time.split(':').fold(0.0, |total, part| { |
| 308 | total * 60.0 + part.parse::<f64>().unwrap_or(0.0) |
| 309 | }); |
| 310 | ( |
| 311 | seconds, |
| 312 | words.next().and_then(|rss| rss.parse().ok()).unwrap_or(0), |
| 313 | ) |
| 314 | } |
| 315 | |
| 316 | /// The relay's `/health`, as numbers by name. |
| 317 | fn health(address: &str) -> HashMap<String, f64> { |
| 318 | let mut stream = TcpStream::connect(address).unwrap(); |
| 319 | write!( |
| 320 | stream, |
| 321 | "GET /health HTTP/1.1\r\nHost: relay\r\nConnection: close\r\n\r\n" |
| 322 | ) |
| 323 | .unwrap(); |
| 324 | let mut text = String::new(); |
| 325 | stream.read_to_string(&mut text).unwrap(); |
| 326 | let body = text.split("\r\n\r\n").nth(1).unwrap_or_default(); |
| 327 | body.trim() |
| 328 | .trim_matches(['{', '}']) |
| 329 | .split(',') |
| 330 | .filter_map(|pair| { |
| 331 | let (name, value) = pair.split_once(':')?; |
| 332 | Some((name.trim_matches('"').to_owned(), value.parse().ok()?)) |
| 333 | }) |
| 334 | .collect() |
| 335 | } |
| 336 | |
| 337 | fn report(what: &str, mut times: Vec<Duration>) { |
| 338 | if times.is_empty() { |
| 339 | println!("{what}: none seen"); |
| 340 | return; |
| 341 | } |
| 342 | times.sort(); |
| 343 | let ms = |at: usize| times[at].as_secs_f64() * 1000.0; |
| 344 | println!( |
| 345 | "{what}: median {:.0} ms, p90 {:.0} ms, p99 {:.0} ms, worst {:.0} ms ({} seen)", |
| 346 | ms(times.len() / 2), |
| 347 | ms(times.len() * 9 / 10), |
| 348 | ms(times.len() * 99 / 100), |
| 349 | ms(times.len() - 1), |
| 350 | times.len() |
| 351 | ); |
| 352 | } |
| 353 | |
| 354 | struct Spawned(Child); |
| 355 | |
| 356 | impl Drop for Spawned { |
| 357 | fn drop(&mut self) { |
| 358 | let _ = self.0.kill(); |
| 359 | let _ = self.0.wait(); |
| 360 | } |
| 361 | } |
| 362 | |
| 363 | fn crowd(binary: &str, peers: usize, typists: usize, window: Duration) { |
| 364 | let port = std::net::TcpListener::bind("127.0.0.1:0") |
| 365 | .unwrap() |
| 366 | .local_addr() |
| 367 | .unwrap() |
| 368 | .port(); |
| 369 | let address = format!("127.0.0.1:{port}"); |
| 370 | // Everyone comes from one address here, which a relay would take for one abuser. |
| 371 | let relay = Spawned( |
| 372 | Command::new(binary) |
| 373 | .args(["--listen", &address]) |
| 374 | .args(["--max-connections-per-address", "1000"]) |
| 375 | .args(["--joins-per-minute", "10000"]) |
| 376 | .args(["--failures-per-minute", "10000"]) |
| 377 | .stdout(Stdio::null()) |
| 378 | .spawn() |
| 379 | .unwrap(), |
| 380 | ); |
| 381 | until("no relay", Duration::from_secs(10), || { |
| 382 | TcpStream::connect(&address).is_ok() |
| 383 | }); |
| 384 | let url = format!("ws://{address}"); |
| 385 | let mut hosting = Spawned( |
| 386 | Command::new(std::env::current_exe().unwrap()) |
| 387 | .args(["host", &url, &typists.max(1).to_string()]) |
| 388 | .stdin(Stdio::piped()) |
| 389 | .stdout(Stdio::piped()) |
| 390 | .spawn() |
| 391 | .unwrap(), |
| 392 | ); |
| 393 | let mut said = BufReader::new(hosting.0.stdout.take().unwrap()).lines(); |
| 394 | let code = said.next().unwrap().unwrap(); |
| 395 | let code = code.strip_prefix("code ").unwrap().to_owned(); |
| 396 | println!("{peers} guests, {typists} typing {KEYS_PER_SECOND} keys a second, code {code}"); |
| 397 | |
| 398 | let directory = tempfile::tempdir().unwrap(); |
| 399 | let began = Instant::now(); |
| 400 | let joining: Vec<_> = (0..peers) |
| 401 | .map(|n| { |
| 402 | let (code, url) = (code.clone(), url.clone()); |
| 403 | let cache = directory.path().join(format!("guest{n}")); |
| 404 | thread::sleep(Duration::from_millis(50)); |
| 405 | thread::spawn(move || join(&format!("Guest {n}"), &code, &url, &cache)) |
| 406 | }) |
| 407 | .collect(); |
| 408 | let guests: Vec<Arc<Open>> = joining |
| 409 | .into_iter() |
| 410 | .map(|joined| Arc::new(joined.join().unwrap())) |
| 411 | .collect(); |
| 412 | println!("all joined in {:.1} s", began.elapsed().as_secs_f64()); |
| 413 | until("not everyone met", Duration::from_secs(300), || { |
| 414 | guests.iter().all(|open| open.guest.peers().len() >= peers) |
| 415 | }); |
| 416 | println!("all met in {:.1} s", began.elapsed().as_secs_f64()); |
| 417 | |
| 418 | let (_, texts) = texts_of(&guests[0].section).unwrap(); |
| 419 | assert!(texts.len() >= typists, "a paragraph for each typist"); |
| 420 | for open in &guests { |
| 421 | open.guest.set_presence(presence(open, 0, 0).unwrap()); |
| 422 | } |
| 423 | |
| 424 | // When each typist typed each key, and said its caret moved past it. |
| 425 | let typed: Arc<Mutex<HashMap<(usize, usize), Instant>>> = Arc::default(); |
| 426 | let moved: Arc<Mutex<HashMap<(usize, usize), Instant>>> = Arc::default(); |
| 427 | let seen_keys: Arc<Mutex<Vec<Duration>>> = Arc::default(); |
| 428 | let seen_carets: Arc<Mutex<Vec<Duration>>> = Arc::default(); |
| 429 | let done = Arc::new(AtomicBool::new(false)); |
| 430 | let (relay_before, host_before, crowd_before) = ( |
| 431 | cost(relay.0.id()), |
| 432 | cost(hosting.0.id()), |
| 433 | cost(std::process::id()), |
| 434 | ); |
| 435 | let bytes_before = health(&address); |
| 436 | let started = Instant::now(); |
| 437 | |
| 438 | let mut threads = Vec::new(); |
| 439 | for (n, open) in guests.iter().enumerate().take(typists) { |
| 440 | let (open, typed, done) = (Arc::clone(open), Arc::clone(&typed), Arc::clone(&done)); |
| 441 | let moved = Arc::clone(&moved); |
| 442 | threads.push(thread::spawn(move || { |
| 443 | let pause = Duration::from_secs_f64(1.0 / KEYS_PER_SECOND); |
| 444 | thread::sleep(pause.mul_f64(n as f64 / typists as f64)); |
| 445 | let mut count = 0; |
| 446 | while !done.load(Ordering::Acquire) && count < 600 { |
| 447 | open.section.events(); |
| 448 | let Some((space, mut texts)) = texts_of(&open.section) else { |
| 449 | continue; |
| 450 | }; |
| 451 | let (text, before) = texts.swap_remove(n); |
| 452 | let end = before.encode_utf16().count() as u32; |
| 453 | let edit = Edit { |
| 454 | at: 134_000_000_000_000_000, |
| 455 | ops: vec![Op::Page { |
| 456 | space, |
| 457 | op: PageOp::Text { |
| 458 | text, |
| 459 | range: end..end, |
| 460 | with: key(n, count).to_string(), |
| 461 | }, |
| 462 | }], |
| 463 | }; |
| 464 | typed.lock().unwrap().insert((n, count), Instant::now()); |
| 465 | open.section.replica().apply("Typist", edit).unwrap(); |
| 466 | count += 1; |
| 467 | if let Some(presence) = presence(&open, 0, count as u32) { |
| 468 | moved.lock().unwrap().insert((n, count), Instant::now()); |
| 469 | open.guest.set_presence(presence); |
| 470 | } |
| 471 | thread::sleep(pause); |
| 472 | } |
| 473 | })); |
| 474 | } |
| 475 | for open in guests.iter().skip(typists).take(OBSERVERS) { |
| 476 | let (open, typed, done) = (Arc::clone(open), Arc::clone(&typed), Arc::clone(&done)); |
| 477 | let moved = Arc::clone(&moved); |
| 478 | let (keys, carets) = (Arc::clone(&seen_keys), Arc::clone(&seen_carets)); |
| 479 | threads.push(thread::spawn(move || { |
| 480 | let mut keys_seen = std::collections::HashSet::new(); |
| 481 | let mut carets_seen = HashMap::new(); |
| 482 | while !done.load(Ordering::Acquire) { |
| 483 | open.section.events(); |
| 484 | let now = Instant::now(); |
| 485 | if let Some((_, texts)) = texts_of(&open.section) { |
| 486 | let chars = texts.iter().flat_map(|(_, text)| text.chars()); |
| 487 | for (typist, n) in chars.filter_map(typist_of) { |
| 488 | if keys_seen.insert((typist, n)) |
| 489 | && let Some(at) = typed.lock().unwrap().get(&(typist, n)) |
| 490 | { |
| 491 | keys.lock().unwrap().push(now - *at); |
| 492 | } |
| 493 | } |
| 494 | } |
| 495 | for peer in open.guest.peers() { |
| 496 | let Some(typist) = (peer.hello.name.strip_prefix("Guest ")) |
| 497 | .and_then(|n| n.parse::<usize>().ok()) |
| 498 | .filter(|n| *n < typists) |
| 499 | else { |
| 500 | continue; |
| 501 | }; |
| 502 | let Some(count) = peer.presence.and_then(|p| p.caret).map(|c| c.focus.offset) |
| 503 | else { |
| 504 | continue; |
| 505 | }; |
| 506 | let count = count as usize; |
| 507 | if count > 0 |
| 508 | && carets_seen.insert(typist, count) != Some(count) |
| 509 | && let Some(at) = moved.lock().unwrap().get(&(typist, count)) |
| 510 | { |
| 511 | carets.lock().unwrap().push(now - *at); |
| 512 | } |
| 513 | } |
| 514 | thread::sleep(LOOK); |
| 515 | } |
| 516 | })); |
| 517 | } |
| 518 | thread::sleep(window); |
| 519 | done.store(true, Ordering::Release); |
| 520 | let elapsed = started.elapsed().as_secs_f64(); |
| 521 | let (relay_after, host_after, crowd_after) = ( |
| 522 | cost(relay.0.id()), |
| 523 | cost(hosting.0.id()), |
| 524 | cost(std::process::id()), |
| 525 | ); |
| 526 | let bytes_after = health(&address); |
| 527 | for thread in threads { |
| 528 | thread.join().unwrap(); |
| 529 | } |
| 530 | let keys = typed.lock().unwrap().len(); |
| 531 | report( |
| 532 | "keys, typist to observer", |
| 533 | std::mem::take(&mut seen_keys.lock().unwrap()), |
| 534 | ); |
| 535 | report( |
| 536 | "carets, typist to observer", |
| 537 | std::mem::take(&mut seen_carets.lock().unwrap()), |
| 538 | ); |
| 539 | let delta = |name: &str| bytes_after[name] - bytes_before.get(name).copied().unwrap_or(0.0); |
| 540 | let members = (peers + 1) as f64; |
| 541 | println!( |
| 542 | "relay passed on {:.0} KB/s in, {:.0} KB/s out: per peer {:.1} KB/s in, {:.1} KB/s out", |
| 543 | delta("bytes_in") / elapsed / 1000.0, |
| 544 | delta("bytes_out") / elapsed / 1000.0, |
| 545 | delta("bytes_in") / elapsed / 1000.0 / members, |
| 546 | delta("bytes_out") / elapsed / 1000.0 / members, |
| 547 | ); |
| 548 | let load = |name: &str, (before, _): (f64, u64), (after, rss): (f64, u64)| { |
| 549 | println!( |
| 550 | "{name}: {:.0}% of a core, {:.1} MB resident", |
| 551 | (after - before) / elapsed * 100.0, |
| 552 | rss as f64 / 1024.0 |
| 553 | ); |
| 554 | }; |
| 555 | load("relay", relay_before, relay_after); |
| 556 | load("host", host_before, host_after); |
| 557 | load(&format!("{peers} guests"), crowd_before, crowd_after); |
| 558 | |
| 559 | // Every key typed should reach the host's file; conflict pages show as extra pages. |
| 560 | thread::sleep(Duration::from_secs(5)); |
| 561 | let observed = texts_of(&guests[typists.min(peers - 1)].section).map_or(0, |(_, texts)| { |
| 562 | let chars = texts.iter().flat_map(|(_, text)| text.chars()); |
| 563 | chars.filter_map(typist_of).count() |
| 564 | }); |
| 565 | println!("{keys} keys typed, {observed} on an observer's page after 5 s"); |
| 566 | drop(hosting.0.stdin.take()); |
| 567 | for line in said.map_while(Result::ok) { |
| 568 | println!("host: {line}"); |
| 569 | } |
| 570 | } |