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
14use notebook::{
15 Replica,
16 live::{
17 Caret, Guid, Hello, Presence, Spot,
18 share::{self, Guest, Host, Sharing},
19 },
20 session::{Background, Notebook, Section},
21};
22use onestore::{
23 ExGuid,
24 op::{Edit, Op, PageOp},
25};
26use 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.
41const KEYS_PER_SECOND: f64 = 5.0;
42const LOOK: Duration = Duration::from_millis(2);
43/// Guests that time what they see.
44const OBSERVERS: usize = 3;
45
46fn 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
54fn hello(name: &str) -> Hello {
55 Hello::new(name.into(), None).unwrap()
56}
57
58fn 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.
82fn 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.
121fn 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.
165struct Open {
166 guest: Arc<Guest>,
167 _background: Background,
168 section: Section,
169 file: [u8; 16],
170}
171
172fn 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.
206fn 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.
231fn 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.
271fn 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
277fn 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.
289fn key(typist: usize, n: usize) -> char {
290 char::from_u32(0x4E00 + (typist * 600 + n) as u32).unwrap()
291}
292
293fn 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`.
299fn 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.
317fn 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
337fn 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
354struct Spawned(Child);
355
356impl Drop for Spawned {
357 fn drop(&mut self) {
358 let _ = self.0.kill();
359 let _ = self.0.wait();
360 }
361}
362
363fn 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}