1//! Replays recorded live answers, keyed by the exact request bytes. 2//! 3//! A fixture is JSON: `{"model", "exchanges": [{"request", "responses": 4//! [{"status", "headers": [[name, value]], "body"}]}]}`, where `request` 5//! and `body` are the bytes on the wire as text. Each request's responses 6//! are served in order; today every list holds exactly one. A request not 7//! in the fixture is answered with a non-retryable 400 naming it, and 8//! [`Replay::assert_complete`] (also run on drop) fails the test with the 9//! request, so a changed request can never reach the network or pass. 10 11use std::collections::HashMap; 12use std::sync::{Arc, Mutex}; 13 14use bytes::Bytes; 15 16use crate::{Recorded, Reply}; 17 18/// 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} 25 26/// A parsed fixture file. 27#[derive(Debug, Clone)] 28pub struct Fixture { 29 pub model: String, 30 pub exchanges: Vec<(String, Vec<Response>)>, 31} 32 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 } 65 66 /// 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} 85 86#[derive(Default)] 87struct State { 88 served: HashMap<String, usize>, 89 misses: Vec<String>, 90} 91 92/// The responder side of a replay; hand [`Replay::responder`] to 93/// [`MockJev::start`](crate::MockJev::start). 94#[derive(Clone)] 95pub struct Replay { 96 by_request: Arc<HashMap<String, Vec<Response>>>, 97 state: Arc<Mutex<State>>, 98} 99 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 } 135 136 /// 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 } 140 141 /// 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} 150 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} 167