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}