1//! `/journal/push` and `/journal/pull`: the one log. A device sends the entries the service has not 2//! received from it and fetches the lines of the log it has not seen; the log only grows, so nothing 3//! can conflict. Parents see every conversation together, wherever it happened. 4 5use ::log::{debug, error, info, warn}; 6use serde_json::{Value, json}; 7use whiskers_ports::{Appended, DeviceCursor, DeviceName, EntryLine, HasJournal, Journal as _, LogCursor}; 8 9use super::{Route, sealed}; 10use crate::Response; 11 12/// `POST /journal/push` and `/journal/pull`. Needs the journal. 13pub struct Journal; 14 15impl sealed::Sealed for Journal {} 16 17impl<C: HasJournal> Route<C> for Journal { 18 const PATHS: &'static [&'static str] = &["/journal/push", "/journal/pull"]; 19 async fn answer(&self, adapter: &C, route: &str, body: &str) -> Response { 20 journal(adapter, route, body).await 21 } 22} 23 24fn reply(status: u16, json: Value) -> Response { 25 Response::json(status, json.to_string()) 26} 27 28fn failed(route: &str, e: &dyn std::fmt::Display) -> Response { 29 error!("{route}: the journal failed: {e}"); 30 reply(500, json!({ "error": e.to_string() })) 31} 32 33async fn journal<C: HasJournal>(adapter: &C, route: &str, body: &str) -> Response { 34 let Ok(v) = serde_json::from_str::<Value>(body) else { 35 warn!("{route}: the body ({} bytes) is not JSON", body.len()); 36 return reply(400, json!({ "error": "bad json" })); 37 }; 38 if route == "/journal/push" { journal_push(adapter, &v).await } else { journal_pull(adapter, &v).await } 39} 40 41async fn journal_push<C: HasJournal>(adapter: &C, v: &Value) -> Response { 42 let Some(device) = v["device"].as_str().and_then(|d| DeviceName::new(d).ok()) else { 43 warn!("journal push refused: the device name is missing or not plain"); 44 return reply(400, json!({ "error": "device" })); 45 }; 46 let from = DeviceCursor::new(v["from"].as_u64().unwrap_or(0)); 47 let raw: Vec<&str> = v["lines"].as_array().map(|a| a.iter().filter_map(|l| l.as_str()).collect()).unwrap_or_default(); 48 let held = match adapter.journal().held(&device).await { 49 Ok(held) => held, 50 Err(e) => return failed("/journal/push", &e), 51 }; 52 // Only lines that continue exactly where the service's copy of this device ends are taken; a 53 // device that was behind or ahead is told the real count and resends from there. (A push that 54 // will not be taken is not inspected: it is answered with the count, whatever it carries.) 55 debug!("journal push from {}: {} line(s) at {}, the service has {}", device.as_str(), raw.len(), from.get(), held.get()); 56 if raw.is_empty() { 57 return reply(200, json!({ "have": held.get() })); 58 } 59 if from != held { 60 warn!("journal push from {} starts at {} but the service has {}; telling it the real count", device.as_str(), from.get(), held.get()); 61 return reply(200, json!({ "have": held.get() })); 62 } 63 let mut lines = Vec::with_capacity(raw.len()); 64 for l in &raw { 65 match EntryLine::parse(l) { 66 Ok(entry) => lines.push(entry), 67 Err(_) => { 68 warn!("journal push from {} refused: a line is not a single JSON object", device.as_str()); 69 return reply(400, json!({ "error": "lines must be single JSON objects" })); 70 } 71 } 72 } 73 match adapter.journal().append(&device, from, &lines).await { 74 Ok(Appended::Accepted { held }) => { 75 info!("journal: {} line(s) appended for {}; it now has {}", lines.len(), device.as_str(), held.get()); 76 reply(200, json!({ "have": held.get() })) 77 } 78 Ok(Appended::Misaligned { held }) => { 79 // Another push from the same device got in between: nothing was taken here. 80 warn!("journal push from {} lost a race: the service has {}; telling it the real count", device.as_str(), held.get()); 81 reply(200, json!({ "have": held.get() })) 82 } 83 Err(e) => failed("/journal/push", &e), 84 } 85} 86 87/// The log from line `have` on. 88async fn journal_pull<C: HasJournal>(adapter: &C, v: &Value) -> Response { 89 let since = LogCursor::new(v["have"].as_u64().unwrap_or(0)); 90 match adapter.journal().pull(since).await { 91 Ok(page) => { 92 debug!("journal pull from {}: {} of {} line(s) returned", page.from.get(), page.lines.len(), page.total.get()); 93 let lines: Vec<&str> = page.lines.iter().map(|l| l.as_str()).collect(); 94 reply(200, json!({ "from": page.from.get(), "total": page.total.get(), "lines": lines })) 95 } 96 Err(e) => failed("/journal/pull", &e), 97 } 98}