whiskers.git / backend / http-conformance / src / concurrency.rs

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.

10use std::collections::BTreeSet;
11use std::sync::Barrier;
13use serde_json::Value;
14
15use crate::host::Host;

Which route the overlapping requests use: the public sync (which carries no changed) or the test mode's window onto the object (which does).

19#[derive(Clone, Copy, PartialEq, Eq, Debug)]
20pub enum Via {
21    Public,
22    Hub,
23}

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.

41pub struct Tally {
42    pub failures: Vec<String>,
43    pub facts: usize,
44    pub distinct_ids: usize,
45    pub claimed_changes: usize,
46}

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}