Requests that overlap, against one household: the atomicity the ports promise, asked of a running host.
Two invariants, both about a merge being atomic. If two merges could interleave, a lost update would
drop a fact from the final document (1), and two merges could each count the same new fact as theirs
(2). So: after THREADS devices each send EACH distinct facts at once, the household holds all
of them under distinct numbers, and the changed counts the replies reported add up to exactly the
number of facts. And when every device sends the same one fact at the same moment, exactly one of
them is told it was new.
Which route the overlapping requests use: the public sync (which carries no changed) or the test mode's
window onto the object (which does).
Requests that had to be sent again (public route only).
26pub static RETRIES: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
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}
What the replies of /memory/sync said, summed: how many facts the final document holds, whether their
numbers are distinct, and the number of changes the replies claimed.
Counts changes from outside: the wire carries no changed, so the claim is how many replies first showed a
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}
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}
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}