main.rsannotatedmain.rssource123 lines · 6.0 KB · raw
1//! Whiskers' service, run on the operator's machine and reached by the tablet over a private network. It
2//! holds the keys so the child's tablet never does:
3//!
4//! - `POST /check`  a message in, a `Verdict` out (Jev).
5//! - `POST /speak`  `{"text": "..."}` in, `audio/mpeg` out (ElevenLabs).
6//! - `POST /speak/v2`  the same, as JSON with the time of every character (Whiskers' own shape).
7//! - `POST /usage`  what the parents' screen shows about the natural voice: credits, allowance, last failure.
8//! - `POST /memory/sync`, `POST /household/sync`, `POST /chat/sync`  the master copies of what Whiskers remembers,
9//!   of the grown-ups' choices and of the one session; a device sends its own and gets the merged result back.
10//! - `POST /journal/push`, `/journal/pull`, `/picture/missing`, `/picture/put`, `/picture/get`  the one log and the pictures.
11//! - `POST /embed`  texts in, vectors out (a local embedding model).
12//! - `POST /rerank`  a message and a shortlist of memories in, a probability each out (Jev).
13//! - `POST /v1/messages`  the model, through the operator's gateway, narrowed to what Whiskers sends.
14//!
15//!     whiskersd <private-address>:47900
16//!
17//! This file is only the HTTP: it turns a `tiny_http` request into a `whiskers_service::Request`, runs
18//! it, and turns the `Response` back. What the routes do is `whiskers-service`; what they run on is the
19//! native adapters in this crate's library.
20//!
21//! Anything it cannot do (no key, upstream down, spend refused) is a refusal the
22//! tablet turns into a spoken fallback: `Unavailable` for `/check`, an error
23//! status for `/speak`.
24
25mod access;
26
27use std::sync::Arc;
28use std::sync::atomic::{AtomicUsize, Ordering};
29use std::thread;
30
31use access::Access;
32use log::{error, info, trace, warn};
33use whiskers_ports::run_ready;
34use whiskers_service::{Body, Complete, Content, MAX_BODY, Method, Request, Response, Service};
35use whiskersd::Native;
36
37/// Requests served at once, past which new ones are refused with 503 (the tablet then falls
38/// back). One child talks to one cat, but the memory work that follows a turn runs alongside the
39/// next turn, and a slow model call must never hold up a check.
40const MAX_IN_FLIGHT: usize = 64;
41
42fn main() {
43    whiskers_core::logging::init("whiskersd");
44    let addr = std::env::args().nth(1).unwrap_or_else(|| {
45        eprintln!("usage: whiskersd <addr:port>");
46        std::process::exit(2);
47    });
48    info!("whiskersd starting, listening address {addr}");
49    let service = Arc::new(Service::complete(Native::from_env().unwrap_or_else(|e| {
50        error!("cannot start: {e}");
51        eprintln!("whiskersd: {e}");
52        std::process::exit(e.exit_code());
53    })));
54    let server = Arc::new(tiny_http::Server::http(&addr).unwrap_or_else(|e| {
55        error!("cannot listen on {addr}: {e}");
56        std::process::exit(1);
57    }));
58    info!("listening on {addr}");
59
60    // One thread per request, not a fixed pool: a request body is read inline, and a tablet that
61    // drops off wifi halfway through a picture would otherwise pin a pooled worker for good, until
62    // enough of them left no one to answer a safety check.
63    let in_flight = Arc::new(AtomicUsize::new(0));
64    while let Ok(request) = server.recv() {
65        let busy = in_flight.fetch_add(1, Ordering::SeqCst);
66        if busy >= MAX_IN_FLIGHT {
67            in_flight.fetch_sub(1, Ordering::SeqCst);
68            Access::new(&request).respond(request, tiny_http::Response::empty(503));
69            warn!("refused a request: {busy} already in flight (limit {MAX_IN_FLIGHT})");
70            continue;
71        }
72        trace!("request accepted; {} in flight", busy + 1);
73        let (service, in_flight) = (service.clone(), in_flight.clone());
74        thread::spawn(move || {
75            handle(&service, request);
76            in_flight.fetch_sub(1, Ordering::SeqCst);
77        });
78    }
79    warn!("the server stopped receiving requests; exiting");
80}
81
82fn header(name: &str, value: &str) -> tiny_http::Header {
83    tiny_http::Header::from_bytes(name, value).expect("a header made of plain text")
84}
85
86/// Reads a request into the service's terms. Only a POST has its body read; anything else is answered
87/// with a 405 without waiting for one.
88fn read(request: &mut tiny_http::Request, access: &mut Access) -> Request {
89    let method = match request.method() {
90        tiny_http::Method::Post => Method::Post,
91        other => Method::Other(other.to_string()),
92    };
93    let path = request.url().to_owned();
94    if method != Method::Post {
95        return Request { method, path, body: Body::Text(String::new()) };
96    }
97    let mut body = String::new();
98    let read = std::io::Read::read_to_string(&mut std::io::Read::take(request.as_reader(), MAX_BODY as u64 + 1), &mut body);
99    access.read(body.len());
100    let body = match read {
101        Ok(_) => Body::Text(body),
102        Err(e) => Body::Unreadable(whiskers_ports::Diagnostic::new(e.to_string())),
103    };
104    Request { method, path, body }
105}
106
107fn respond(access: &Access, request: tiny_http::Request, response: Response) {
108    let Response { status, content, body } = response;
109    match content {
110        Content::None => access.respond(request, tiny_http::Response::empty(status)),
111        Content::Text => access.respond(request, tiny_http::Response::from_string(String::from_utf8_lossy(&body).into_owned()).with_status_code(status)),
112        Content::Json => access.respond(request, tiny_http::Response::from_data(body).with_status_code(status).with_header(header("Content-Type", "application/json"))),
113        Content::Mpeg => access.respond(request, tiny_http::Response::from_data(body).with_status_code(status).with_header(header("content-type", "audio/mpeg"))),
114    }
115}
116
117fn handle(service: &Service<Native, Complete>, mut request: tiny_http::Request) {
118    let mut access = Access::new(&request);
119    let parsed = read(&mut request, &mut access);
120    // The native adapters block inside their `async fn`s, so every future here is ready on its first poll.
121    let response = run_ready(service.handle(parsed));
122    respond(&access, request, response);
123}