whiskers.git / backend / http-conformance / src / concurrency.rs
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}