whiskers.git / crates / whiskersd / src / judge.rs
judge.rsannotatedjudge.rssource147 lines · 6.1 KB · raw

Jev, as the judge. The client lives on its own thread (its futures are not Send) and the request threads hand it jobs; one that hangs must not hold every request behind it, so a request waits a bounded time for an answer. There is one judge per supported age, built at start: the questions are worded for the age they judge.

6use std::sync::Arc;
7use std::sync::mpsc::{Sender, channel};
8use std::thread;
10use jev_http::{Endpoint, Jev, Ledger};
11use jev_protocol::ModelId;
12use log::{debug, error, info, warn};
13use whiskers_core::{Age, Direction, Verdict};
14use whiskers_guard::{Judge as Questions, MODEL, Rerank};
15use whiskers_ports::{Diagnostic, Judge, JudgeError, Lookup, SecretName, Secrets, run_ready};

How long a request waits for the Jev thread before giving up (the tablet then falls back).

19const JEV_PATIENCE: std::time::Duration = std::time::Duration::from_secs(20);

A job for the thread that owns the Jev client.

22type Job = Box<dyn FnOnce(&Jev, &tokio::runtime::Runtime) + Send>;
24pub struct JevJudge {
25    judges: Vec<(Age, Arc<Questions>)>,

None when there is no key: every question is then NotConfigured.

27    jev: Option<Sender<Job>>,
28}
30impl JevJudge {

Builds the judges and, if TYPESAFE_API_KEY is set, the client and its thread.

32    pub fn start(secrets: &impl Secrets) -> Self {
33        let judges = (Age::YOUNGEST.years()..=Age::OLDEST.years())
34            .map(|y| {
35                let age = Age::new(y).expect("within the supported range");
36                (age, Arc::new(Questions::new(age).expect("the guard questions are valid")))
37            })
38            .collect();
39        Self { judges, jev: start_jev(secrets) }
40    }

Runs f on the Jev thread and waits for its answer.

43    fn with_jev<T: Send + 'static>(&self, f: impl FnOnce(&Jev, &tokio::runtime::Runtime) -> T + Send + 'static) -> Result<T, JudgeError> {
44        let Some(tx) = self.jev.as_ref() else {
45            warn!("a Jev question was asked but there is no Jev client (no key)");
46            return Err(JudgeError::NotConfigured(SecretName::JevKey));
47        };
48        let started = std::time::Instant::now();
49        let (reply_tx, reply_rx) = channel();
50        let job: Job = Box::new(move |jev, rt| {
51            let _ = reply_tx.send(f(jev, rt));
52        });
53        if tx.send(job).is_err() {
54            error!("the Jev thread is gone; no question can be asked");
55            return Err(JudgeError::NoAnswer);
56        }
57        match reply_rx.recv_timeout(JEV_PATIENCE) {
58            Ok(v) => {
59                debug!("Jev thread answered in {} ms", started.elapsed().as_millis());
60                Ok(v)
61            }
62            Err(e) => {
63                error!("no answer from the Jev thread after {} ms: {e}", started.elapsed().as_millis());
64                Err(JudgeError::NoAnswer)
65            }
66        }
67    }
68}
70fn start_jev(secrets: &impl Secrets) -> Option<Sender<Job>> {
71    let key = match run_ready(secrets.get(SecretName::JevKey)) {
72        Ok(Lookup::Present(key)) => key,
73        Ok(Lookup::Empty) => {
74            warn!("TYPESAFE_API_KEY is set but empty; every /check and /rerank will be Unavailable");
75            return None;
76        }
77        Ok(Lookup::Absent) => {
78            warn!("TYPESAFE_API_KEY is not set; every /check and /rerank will be Unavailable");
79            return None;
80        }
81        Err(e) => {
82            error!("TYPESAFE_API_KEY could not be looked up: {e}");
83            return None;
84        }
85    };
86    let jev = match Ledger::open()
87        .map_err(|e| format!("opening the spend ledger: {e}"))
88        .and_then(|ledger| Jev::new(key.expose(), "whiskersd", ledger, ModelId::pinned(MODEL).expect("model id"), Endpoint::api()))
89    {
90        Ok(jev) => jev,
91        Err(e) => {
92            error!("cannot build the Jev client: {e}");
93            return None;
94        }
95    };
96    info!("Jev client ready, model {MODEL}");
97    let (tx, rx) = channel::<Job>();
98    thread::spawn(move || {
99        let rt = tokio::runtime::Builder::new_current_thread().enable_all().build().expect("runtime");
100        for job in rx {
101            // A job that panics must not end the only thread that can talk to Jev.
102            if std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| job(&jev, &rt))).is_err() {
103                error!("a Jev job panicked; the Jev thread carries on");
104            }
105        }
106    });
107    Some(tx)
108}
109
110impl Judge for JevJudge {
111    async fn check(&self, direction: Direction, age: Age, text: &str) -> Result<Verdict, JudgeError> {
112        let Some((_, judge)) = self.judges.iter().find(|(a, _)| *a == age) else {
113            error!("/check: no judge for age {}", age.years());
114            return Err(JudgeError::NoJudgeFor(age));
115        };
116        let judge = judge.clone();
117        let text = text.to_owned();
118        self.with_jev(move |jev, rt| {
119            let js = Questions::state(direction, age, &text);
120            match rt.block_on(jev.ask(&js, &judge.questions, &format!("whiskers: is this message suitable for a {}-year-old", age.years()))) {
121                Ok(answered) => Ok(judge.verdict(direction, &answered.response)),
122                Err(e) => {
123                    error!("Jev failed on /check: {e}");
124                    Err(JudgeError::Failed(Diagnostic::new(e.to_string())))
125                }
126            }
127        })?
128    }
129
130    async fn rerank(&self, query: &str, candidates: &[String]) -> Result<Vec<f32>, JudgeError> {
131        let ranking = Rerank::new(candidates).map_err(|e| {
132            warn!("/rerank: cannot rank {} candidates: {e}", candidates.len());
133            JudgeError::CannotRank { candidates: candidates.len(), why: Diagnostic::new(e.to_string()) }
134        })?;
135        let query = query.to_owned();
136        self.with_jev(move |jev, rt| {
137            let js = Rerank::state(&query);
138            match rt.block_on(jev.ask(&js, &ranking.questions, "whiskers: which memories fit what the child said")) {
139                Ok(answered) => Ok(ranking.probabilities(&answered.response)),
140                Err(e) => {
141                    error!("Jev failed on /rerank: {e}");
142                    Err(JudgeError::Failed(Diagnostic::new(e.to_string())))
143                }
144            }
145        })?
146    }
147}