1//! Requests that overlap, against one household: the atomicity the ports promise, asked of a running host. 2//! 3//! Two invariants, both about a merge being atomic. If two merges could interleave, a lost update would 4//! drop a fact from the final document (1), and two merges could each count the same new fact as theirs 5//! (2). So: after `THREADS` devices each send `EACH` distinct facts at once, the household holds all 6//! of them under distinct numbers, and the `changed` counts the replies reported add up to exactly the 7//! number of facts. And when every device sends the same one fact at the same moment, exactly one of 8//! them is told it was new. 9 10use std::collections::BTreeSet; 11use std::sync::Barrier; 12 13use serde_json::Value; 14 15use crate::host::Host; 16 17/// Which route the overlapping requests use: the public sync (which carries no `changed`) or the test mode's 18/// window onto the object (which does). 19#[derive(Clone, Copy, PartialEq, Eq, Debug)] 20pub enum Via { 21 Public, 22 Hub, 23} 24 25/// Requests that had to be sent again (public route only). 26pub static RETRIES: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0); 27 28pub const THREADS: usize = 16; 29pub const EACH: usize = 8; 30 31fn fact(gid: &str) -> String { 32 format!(r#"{{"id":0,"text":"fact number {gid}","learned_at_ms":5,"kind":"Other","who":[],"place":null,"when":null,"pictures":[],"embedding":[],"gid":"{gid}"}}"#) 33} 34 35fn document(reply: &str) -> Result<Value, String> { 36 serde_json::from_str(reply).map_err(|e| format!("a reply that is not JSON: {e}")) 37} 38 39/// What the replies of `/memory/sync` said, summed: how many facts the final document holds, whether their 40/// numbers are distinct, and the number of changes the replies claimed. 41pub struct Tally { 42 pub failures: Vec<String>, 43 pub facts: usize, 44 pub distinct_ids: usize, 45 pub claimed_changes: usize, 46} 47 48/// Counts changes from outside: the wire carries no `changed`, so the claim is how many replies first showed a 49/// fact. Instead the test-mode `/_hub/memory/merge` route is asked, which does carry it. 50fn merge(host: &Host, household: &str, body: &str, via: Via) -> Result<(Value, usize), String> { 51 let path = match via { 52 Via::Hub => "/_hub/memory/merge", 53 Via::Public => "/memory/sync", 54 }; 55 // A device whose sync fails tries again, and a sync is idempotent, so the public route is retried (up to 3 56 // times). The count is reported: `wrangler dev` sometimes loses a connection in its own proxy. 57 let attempts = if via == Via::Public { 3 } else { 1 }; 58 let mut reply = host.post(path, household, body)?; 59 for _ in 1..attempts { 60 if reply.status == 200 { 61 break; 62 } 63 RETRIES.fetch_add(1, std::sync::atomic::Ordering::Relaxed); 64 reply = host.post(path, household, body)?; 65 } 66 if reply.status != 200 { 67 return Err(format!("status {}", reply.status)); 68 } 69 let v = document(&reply.body)?; 70 if via == Via::Public { 71 return Ok((v, 0)); 72 } 73 let changed = v["changed"].as_u64().ok_or("no changed")? as usize; 74 Ok((v["document"].clone(), changed)) 75} 76 77pub fn distinct_facts(host: &Host, household: &str, via: Via) -> Tally { 78 let barrier = Barrier::new(THREADS); 79 let results: Vec<Vec<Result<(Value, usize), String>>> = std::thread::scope(|s| { 80 let handles: Vec<_> = (0..THREADS) 81 .map(|t| { 82 let barrier = &barrier; 83 s.spawn(move || { 84 barrier.wait(); 85 (0..EACH) 86 .map(|i| merge(host, household, &format!(r#"{{"facts":[{}],"forgotten":[]}}"#, fact(&format!("t{t:02}f{i:02}"))), via)) 87 .collect() 88 }) 89 }) 90 .collect(); 91 handles.into_iter().map(|h| h.join().expect("a worker thread")).collect() 92 }); 93 let mut failures = Vec::new(); 94 let mut claimed = 0; 95 for r in results.into_iter().flatten() { 96 match r { 97 Ok((_, c)) => claimed += c, 98 Err(e) => failures.push(e), 99 } 100 } 101 let (facts, distinct) = match merge(host, household, r#"{"facts":[],"forgotten":[]}"#, via) { 102 Ok((doc, _)) => { 103 let facts = doc["facts"].as_array().cloned().unwrap_or_default(); 104 let ids: BTreeSet<u64> = facts.iter().filter_map(|f| f["id"].as_u64()).collect(); 105 (facts.len(), ids.len()) 106 } 107 Err(e) => { 108 failures.push(format!("the final read: {e}")); 109 (0, 0) 110 } 111 }; 112 Tally { failures, facts, distinct_ids: distinct, claimed_changes: claimed } 113} 114 115/// Everyone sends the same fact at once: returns how many replies claimed it as new. 116pub fn the_same_fact(host: &Host, household: &str) -> Result<usize, String> { 117 let via = Via::Hub; 118 let barrier = Barrier::new(THREADS); 119 let body = format!(r#"{{"facts":[{}],"forgotten":[]}}"#, fact("same")); 120 let results: Vec<Result<(Value, usize), String>> = std::thread::scope(|s| { 121 let handles: Vec<_> = (0..THREADS) 122 .map(|_| { 123 let (barrier, body) = (&barrier, &body); 124 s.spawn(move || { 125 barrier.wait(); 126 merge(host, household, body, via) 127 }) 128 }) 129 .collect(); 130 handles.into_iter().map(|h| h.join().expect("a worker thread")).collect() 131 }); 132 let mut new = 0; 133 for r in results { 134 new += r?.1; 135 } 136 Ok(new) 137}