Replays recorded live answers, keyed by the exact request bytes.

A fixture is JSON: {"model", "exchanges": [{"request", "responses": [{"status", "headers": [[name, value]], "body"}]}]}, where request and body are the bytes on the wire as text. Each request's responses are served in order; today every list holds exactly one. A request not in the fixture is answered with a non-retryable 400 naming it, and [Replay::assert_complete] (also run on drop) fails the test with the request, so a changed request can never reach the network or pass.

11use std::collections::HashMap;
12use std::sync::{Arc, Mutex};
14use bytes::Bytes;
15
16use crate::{Recorded, Reply};

One recorded answer.

19#[derive(Debug, Clone)]
20pub struct Response {
21    pub status: u16,
22    pub headers: Vec<(String, String)>,
23    pub body: String,
24}

A parsed fixture file.

27#[derive(Debug, Clone)]
28pub struct Fixture {
29    pub model: String,
30    pub exchanges: Vec<(String, Vec<Response>)>,
31}
33impl Fixture {
34    pub fn parse(text: &str) -> Result<Self, String> {
35        let v: serde_json::Value = serde_json::from_str(text).map_err(|e| e.to_string())?;
36        let str_of = |v: &serde_json::Value, what: &str| {
37            v.as_str().map(str::to_owned).ok_or_else(|| format!("fixture: {what} is not a string"))
38        };
39        let model = str_of(&v["model"], "model")?;
40        let mut exchanges = Vec::new();
41        for (i, e) in v["exchanges"].as_array().ok_or("fixture: no exchanges")?.iter().enumerate() {
42            let request = str_of(&e["request"], &format!("exchanges[{i}].request"))?;
43            let mut responses = Vec::new();
44            for r in e["responses"].as_array().ok_or(format!("fixture: exchanges[{i}] has no responses"))? {
45                let status = r["status"].as_u64().and_then(|s| u16::try_from(s).ok()).ok_or("fixture: bad status")?;
46                let mut headers = Vec::new();
47                for h in r["headers"].as_array().map(Vec::as_slice).unwrap_or_default() {
48                    headers.push((str_of(&h[0], "header name")?, str_of(&h[1], "header value")?));
49                }
50                responses.push(Response { status, headers, body: str_of(&r["body"], "body")? });
51            }
52            if responses.is_empty() {
53                return Err(format!("fixture: exchanges[{i}] has no responses"));
54            }
55            exchanges.push((request, responses));
56        }
57        Ok(Fixture { model, exchanges })
58    }
59
60    pub fn load(path: impl AsRef<std::path::Path>) -> Result<Self, String> {
61        let path = path.as_ref();
62        let text = std::fs::read_to_string(path).map_err(|e| format!("{}: {e}", path.display()))?;
63        Self::parse(&text).map_err(|e| format!("{}: {e}", path.display()))
64    }

The fixture as the recorder writes it.

67    pub fn to_json(&self) -> String {
68        let exchanges: Vec<_> = self
69            .exchanges
70            .iter()
71            .map(|(request, responses)| {
72                serde_json::json!({
73                    "request": request,
74                    "responses": responses.iter().map(|r| serde_json::json!({
75                        "status": r.status,
76                        "headers": r.headers.iter().map(|(n, v)| [n, v]).collect::<Vec<_>>(),
77                        "body": r.body,
78                    })).collect::<Vec<_>>(),
79                })
80            })
81            .collect();
82        serde_json::to_string_pretty(&serde_json::json!({ "model": self.model, "exchanges": exchanges })).unwrap()
83    }
84}
86#[derive(Default)]
87struct State {
88    served: HashMap<String, usize>,
89    misses: Vec<String>,
90}

The responder side of a replay; hand [Replay::responder] to MockJev::start.

94#[derive(Clone)]
95pub struct Replay {
96    by_request: Arc<HashMap<String, Vec<Response>>>,
97    state: Arc<Mutex<State>>,
98}
100impl Replay {
101    pub fn new(fixture: &Fixture) -> Self {
102        Replay {
103            by_request: Arc::new(fixture.exchanges.iter().cloned().collect()),
104            state: Arc::default(),
105        }
106    }
107
108    pub fn responder(&self) -> impl Fn(&Recorded) -> Reply + Send + Sync + 'static {
109        let this = self.clone();
110        move |req| this.answer(req)
111    }
112
113    fn answer(&self, req: &Recorded) -> Reply {
114        let request = String::from_utf8_lossy(&req.body).into_owned();
115        let mut state = self.state.lock().unwrap();
116        let next = state.served.get(&request).copied().unwrap_or(0);
117        match self.by_request.get(&request).and_then(|rs| rs.get(next)) {
118            Some(r) => {
119                state.served.insert(request, next + 1);
120                let mut reply = Reply::json(r.status, serde_json::Value::Null);
121                reply.raw_body = Some(Bytes::from(r.body.clone()));
122                for (name, value) in &r.headers {
123                    // Header names are few and live for the test.
124                    reply.headers.push((Box::leak(name.clone().into_boxed_str()), value.clone()));
125                }
126                reply
127            }
128            None => {
129                let message = miss_message(&request, next);
130                state.misses.push(request);
131                Reply::json(400, serde_json::json!({ "detail": { "error_type": "fixture_miss", "message": message } }))
132            }
133        }
134    }

The requests that were not in the fixture, as sent.

137    pub fn misses(&self) -> Vec<String> {
138        self.state.lock().unwrap().misses.clone()
139    }

Panics naming every request the fixture did not hold.

142    pub fn assert_complete(&self) {
143        let misses = self.misses();
144        if !misses.is_empty() {
145            let named: Vec<_> = misses.iter().map(|m| miss_message(m, 0)).collect();
146            panic!("{} request(s) not in the fixture:\n{}", misses.len(), named.join("\n"));
147        }
148    }
149}
151fn miss_message(request: &str, served: usize) -> String {
152    if served > 0 {
153        format!("fixture has no response #{} for request {request}", served + 1)
154    } else {
155        format!("request not in fixture: {request}")
156    }
157}
158
159impl Drop for Replay {
160    fn drop(&mut self) {
161        // The last handle checks, unless the test is already failing.
162        if Arc::strong_count(&self.state) == 1 && !std::thread::panicking() {
163            self.assert_complete();
164        }
165    }
166}