lib.rsannotatedlib.rssource446 lines · 20.0 KB · raw
1//! The Jev API over HTTPS for a native process: `jev-client`'s policy on
2//! this crate's transport (`transport.rs`: HTTP/2 over pure-Rust TLS on
3//! tokio) and runtime, and the shared spend ledger in front of both.
4//!
5//! WHAT to ask is `jev-protocol`'s; how to retry is `jev-client`'s. This is
6//! only how the bytes get there, and what they cost. The browser build cannot
7//! reach the API directly anyway (`api.typesafe.ai` refuses browser
8//! origins), so a web page gets a relay and never links this.
9//!
10//! **Every request is checked against [`ledger::Ledger`] before it is sent.**
11//! This account's credit is finite, so a client that could be built and used
12//! without a shared, on-disk record of what has already been spent would make
13//! overspending a matter of which binary you happened to run - see `ledger`'s
14//! own doc comment for why the record has to be a file and not a number kept
15//! in memory. `Jev::new`/`Jev::from_env` are the only ways to get a client,
16//! and both build the ledger check into it, so there is no path to `ask` that
17//! skips it. There is no self-imposed lifetime cap here (removed 2026-09-22,
18//! the user's own ruling) - the only hard stop is the vendor's own
19//! out-of-credit answer, which is a fact about the account, not a guess.
20
21pub mod ledger;
22mod runtime;
23mod transport;
24
25use std::sync::{Arc, Mutex, PoisonError};
26use std::time::{Duration, Instant};
27
28use jev_client::{Client, ClientError};
29use jev_protocol::{ApiErrorKind, DOLLARS_PER_MTOK, Json, ModelId, Questions, Response, request_bytes, retry};
30
31pub use ledger::{Entry, Guards, Hold, Ledger, NoCredit, Refusal, Status, Window};
32pub use runtime::TokioRuntime;
33pub use transport::{ENDPOINT, Endpoint, HttpsTransport};
34
35/// The environment variable the key arrives in. On this machine it reaches a
36/// process through `op-env-run` and nowhere else - never a dotfile, never a
37/// command line (see the repository README).
38pub const KEY: &str = "TYPESAFE_API_KEY";
39
40/// The documented per-request input ceiling (digest §4), so a huge `state`
41/// cannot make a worst case run away.
42const MAX_REQUEST_TOKENS: f64 = 64_000.0;
43
44/// What a response's headers said about how much of the vendor's own rate
45/// limit is left, when they said anything at all. Digest §10: as of
46/// 2026-09-20 the vendor documents no such headers on a successful response -
47/// only `Retry-After`/`retry-after-ms` on a `429`/`529` - so every field here
48/// is commonly `None`. It stays a typed, if usually-empty, snapshot rather
49/// than nothing at all so that the day the vendor adds one, this starts
50/// reporting it with no code changed at every call site.
51#[derive(
52    Debug, Clone, Copy, Default, PartialEq, serde::Serialize, serde::Deserialize, schemars::JsonSchema,
53)]
54pub struct RateLimit {
55    pub limit: Option<u64>,
56    pub remaining: Option<u64>,
57    /// Seconds until the window resets, exactly as the header gave it - not
58    /// resolved to a clock time, since that would claim a precision about
59    /// when the header was actually read that nothing here can back up.
60    pub reset_seconds: Option<f64>,
61}
62
63/// A client, and one to reuse: it holds the connection, and a warm
64/// connection is worth hundreds of milliseconds per decision.
65pub struct Jev {
66    client: Client<HttpsTransport, TokioRuntime>,
67    /// Which binary this is, for the ledger's records - `native`, `player`,
68    /// `jevprobe`, ...
69    binary: String,
70    ledger: Ledger,
71    guards: Guards,
72    /// The most recent parsed [`RateLimit`], kept for [`Jev::last_rate_limit`]
73    /// so a caller (the `budget` MCP tool, in particular) does not have to
74    /// wait for the next request to see what the last one learned.
75    last_rate_limit: Mutex<Option<RateLimit>>,
76}
77
78/// What one answered request cost in time as well as tokens.
79#[derive(Debug)]
80pub struct Answered {
81    /// The answers, verified against the questions asked.
82    pub response: Response,
83    /// The whole round trip, retries and waits included.
84    pub took: Duration,
85    /// How many times the request had to be sent. 1 is the happy path.
86    pub attempts: u32,
87    /// The response exactly as it arrived. Kept so that a field the protocol
88    /// does not know about is visible rather than silently dropped by the
89    /// parse - which is how the digest gets corrected.
90    pub body: String,
91    /// Every response header, name and value, with anything that looks like
92    /// a credential redacted. `apps/jevprobe --headers` exists to print this.
93    pub headers: Vec<(String, String)>,
94    /// `x-typesafe-request-id`, if the response carried one (digest §1).
95    pub request_id: Option<String>,
96    pub rate_limit: Option<RateLimit>,
97}
98
99/// Why an ask produced no answer. The first two never touch the network.
100#[derive(Debug)]
101pub enum Failure {
102    /// A shared guard that clears itself refused: the same request would be
103    /// admitted after `retry_in`. The caller goes on without an answer.
104    Throttled { why: String, retry_in: Duration },
105    /// Nothing is admitted until a human acts: the vendor said the account is
106    /// out of credit, or the ledger cannot be read or written (a guard that
107    /// cannot see the spend fails closed).
108    Budget(String),
109    /// The request went out and failed, as `jev-client` judged it after its
110    /// retries.
111    Client(ClientError),
112}
113
114impl std::fmt::Display for Failure {
115    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
116        match self {
117            Self::Throttled { why, retry_in } => {
118                write!(f, "not sent, throttled: {why}; clears in {:.0}s", retry_in.as_secs_f64())
119            }
120            Self::Budget(why) => write!(f, "not sent: {why}"),
121            Self::Client(e) => e.fmt(f),
122        }
123    }
124}
125
126impl std::error::Error for Failure {}
127
128impl Jev {
129    /// Read the key from the environment and build a client for the API, its
130    /// ledger opened from `$XDG_STATE_HOME` (see `ledger::Ledger::open`).
131    /// `None` when the key is not there, which is not an error: the app runs
132    /// without it and says on screen that it cannot play.
133    pub fn from_env(binary: impl Into<String>, model: ModelId) -> Option<Result<Self, String>> {
134        let key = std::env::var(KEY).ok().filter(|key| !key.trim().is_empty())?;
135        Some(
136            Ledger::open()
137                .map_err(|e| format!("opening the spend ledger: {e}"))
138                .and_then(|ledger| Self::new(key.trim(), binary, ledger, model, Endpoint::api())),
139        )
140    }
141
142    /// Build a client against an already-open [`Ledger`]. Every real caller
143    /// wants [`Self::from_env`]; this is the constructor that makes the
144    /// dependency explicit rather than reached for internally, so a test can
145    /// hand it a ledger rooted in a private directory instead of the real
146    /// `$XDG_STATE_HOME/jev` - and so there is no `Jev` value anywhere that
147    /// was not required, at construction, to say where its spend is kept.
148    pub fn new(
149        key: &str,
150        binary: impl Into<String>,
151        ledger: Ledger,
152        model: ModelId,
153        endpoint: Endpoint,
154    ) -> Result<Self, String> {
155        // The trust anchors are compiled in rather than read from the host, so
156        // what this trusts is the same on every machine and in a sandbox with
157        // no /etc/ssl at all. The provider is RustCrypto: pure Rust, no ring,
158        // no aws-lc, nothing to cross-compile for the web build later.
159        let mut roots = rustls::RootCertStore::empty();
160        roots.extend(webpki_roots::TLS_SERVER_ROOTS.iter().cloned());
161        if let Some(path) = endpoint.ca_file() {
162            use rustls::pki_types::CertificateDer;
163            use rustls::pki_types::pem::PemObject;
164            let shown = path.display();
165            for cert in CertificateDer::pem_file_iter(path).map_err(|e| format!("CA file {shown}: {e}"))? {
166                let cert = cert.map_err(|e| format!("CA file {shown}: {e}"))?;
167                roots.add(cert).map_err(|e| format!("CA file {shown}: {e}"))?;
168            }
169        }
170        let mut tls = rustls::ClientConfig::builder_with_provider(Arc::new(rustls_rustcrypto::provider()))
171            .with_safe_default_protocol_versions()
172            .map_err(|e| format!("rustls protocol versions: {e}"))?
173            .with_root_certificates(roots)
174            .with_no_client_auth();
175        tls.alpn_protocols = vec![b"h2".to_vec()];
176        let transport = HttpsTransport::new(endpoint, Arc::new(tls));
177        let client = Client::new(transport, TokioRuntime, model, key).map_err(|e| e.to_string())?;
178        Ok(Self {
179            client,
180            binary: binary.into(),
181            ledger,
182            guards: Guards::from_env(),
183            last_rate_limit: Mutex::new(None),
184        })
185    }
186
187    /// The same client under other guards than the environment's.
188    #[must_use]
189    pub fn with_guards(mut self, guards: Guards) -> Self {
190        self.guards = guards;
191        self
192    }
193
194    pub fn model(&self) -> &ModelId {
195        self.client.model()
196    }
197
198    pub fn ledger(&self) -> &Ledger {
199        &self.ledger
200    }
201
202    pub fn guards(&self) -> Guards {
203        self.guards
204    }
205
206    /// What the shared guards say right now - every process's questions, not
207    /// only this one's.
208    pub fn status(&self) -> Result<Status, String> {
209        self.ledger.status(ledger::now(), self.guards)
210    }
211
212    /// The most recent [`RateLimit`] this client has seen, if any response
213    /// ever carried one.
214    pub fn last_rate_limit(&self) -> Option<RateLimit> {
215        *self.last_rate_limit.lock().unwrap_or_else(PoisonError::into_inner)
216    }
217
218    /// Ask, under `jev-client`'s retry policy - but first, have the shared
219    /// ledger admit it ([`Ledger::admit`]: question bucket, hour breaker, out
220    /// of credit - there is no lifetime cap, per the user's own ruling,
221    /// 2026-09-22), before a single byte goes over the network. A refusal
222    /// that clears itself is [`Failure::Throttled`], one that does not (out
223    /// of credit) is [`Failure::Budget`]; neither is retried here.
224    ///
225    /// `why` is recorded against the spend in the ledger's `reason`. It is a
226    /// required argument, not an option, so that every line of the ledger can
227    /// say why the money went: the breaker trip of 2026-09-20 had to be
228    /// reconstructed from token counts because no line said.
229    pub async fn ask(&self, state: &Json, questions: &Questions, why: &str) -> Result<Answered, Failure> {
230        let worst_case = worst_case_usd(self.client.model(), state, questions);
231        let hold = self.ledger.admit(ledger::now(), &self.guards, worst_case, &self.binary, why).map_err(
232            |refusal| match refusal {
233                Refusal::Throttled { why, retry_in } => Failure::Throttled {
234                    why,
235                    retry_in: Duration::try_from_secs_f64(retry_in).unwrap_or(Duration::MAX),
236                },
237                other => Failure::Budget(other.to_string()),
238            },
239        )?;
240
241        let started = Instant::now();
242        match self.client.ask(state, questions).await {
243            Ok(answered) => {
244                let rate_limit = parse_rate_limit(&answered.headers);
245                if let Some(rate_limit) = rate_limit {
246                    *self.last_rate_limit.lock().unwrap_or_else(PoisonError::into_inner) = Some(rate_limit);
247                    if let Err(e) = self.ledger.note_rate_limit(&rate_limit) {
248                        eprintln!("jev ledger: could not record a rate-limit snapshot: {e}");
249                    }
250                }
251                let usage = answered.response.usage();
252                let cost = usage.dollars();
253                let entry = Entry {
254                    time: ledger::now(),
255                    binary: self.binary.clone(),
256                    input_tokens: usage.input_tokens,
257                    output_tokens: usage.output_tokens,
258                    cost_usd: cost,
259                    request_id: answered.request_id.clone(),
260                    estimated: false,
261                    reason: Some(why.to_owned()),
262                    hold: None,
263                };
264                if let Err(e) = self.ledger.settle(&hold, entry) {
265                    // The request already happened and was paid for
266                    // whether or not this write succeeds; refusing to
267                    // return the answer would not undo that, it would
268                    // just also lose the decision. The hold stays
269                    // unsettled, so the ledger goes on counting its
270                    // worst case - more than it cost. Loudly logged,
271                    // since a ledger that cannot be written to is a
272                    // problem worth a human's attention on its own.
273                    eprintln!("jev ledger: could not record a ${cost:.6} spend: {e}");
274                }
275                Ok(Answered {
276                    took: started.elapsed(),
277                    attempts: answered.attempts,
278                    body: String::from_utf8_lossy(&answered.body).into_owned(),
279                    headers: collect_headers(&answered.headers),
280                    request_id: answered.request_id,
281                    rate_limit,
282                    response: answered.response,
283                })
284            }
285            Err(error) => {
286                if let ClientError::Api { error: api, .. } = &error
287                    && api.kind == ApiErrorKind::OutOfCredit
288                    && let Err(e) = self.ledger.mark_out_of_credit("the vendor reports no credit remaining")
289                {
290                    eprintln!("jev ledger: could not record the vendor's out-of-credit answer: {e}");
291                }
292                // `release` books the hold at nothing, which is true only
293                // when nothing can have been billed: the request was never
294                // built, or its one attempt was refused outright. After an
295                // interrupted or timed-out attempt the cost is unknown (the
296                // vendor reports usage only on success), so the hold stays
297                // and goes on counting its worst case, as a crashed
298                // process's does.
299                if nothing_billed(&error)
300                    && let Err(e) = self.ledger.release(&hold, ledger::now(), &error.to_string())
301                {
302                    eprintln!("jev ledger: could not release an unanswered hold: {e}");
303                }
304                Err(Failure::Client(error))
305            }
306        }
307    }
308}
309
310fn nothing_billed(error: &ClientError) -> bool {
311    match error {
312        ClientError::Invalid(_) => true,
313        ClientError::Api { attempts, .. } => *attempts == 1,
314        _ => false,
315    }
316}
317
318/// A conservative worst case, in tokens, for what one attempt could cost -
319/// used to check the guards BEFORE sending, when the true count (which only
320/// the response carries) is not yet known.
321///
322/// A UTF-8 character is never shorter than one token can be cheaper than, so
323/// counting characters can only overstate the token count - which is exactly
324/// what a worst case should do. The documented 64k-token ceiling per request
325/// caps it either way.
326fn worst_case_tokens(model: &ModelId, state: &Json, questions: &Questions) -> f64 {
327    let chars = request_bytes(model, state, questions)
328        .map(|bytes| String::from_utf8_lossy(&bytes).chars().count() as f64)
329        .unwrap_or(MAX_REQUEST_TOKENS);
330    chars.min(MAX_REQUEST_TOKENS)
331}
332
333/// Every attempt `jev-client` may make (`1 + retry::MAX_RETRIES`) billed in
334/// full.
335fn worst_case_usd(model: &ModelId, state: &Json, questions: &Questions) -> f64 {
336    let attempts = f64::from(1 + retry::MAX_RETRIES);
337    worst_case_tokens(model, state, questions) * attempts / 1e6 * DOLLARS_PER_MTOK
338}
339
340/// Any of the conventional rate-limit header spellings, read opportunistically.
341/// Digest §10: none of these are documented anywhere in the vendor's own
342/// reference as of 2026-09-20, so this is expected to return `None` on every
343/// real response; it stays general rather than named after one convention so
344/// that if the vendor starts sending any of them, this starts reporting it
345/// with no code changed anywhere that calls `ask`.
346fn parse_rate_limit(headers: &http::HeaderMap) -> Option<RateLimit> {
347    let text = |name: &str| headers.get(name)?.to_str().ok().map(str::to_owned);
348    let as_u64 = |name: &str| text(name)?.trim().parse::<u64>().ok();
349    let as_f64 = |name: &str| text(name)?.trim().parse::<f64>().ok();
350    let limit = as_u64("x-ratelimit-limit-requests")
351        .or_else(|| as_u64("x-ratelimit-limit"))
352        .or_else(|| as_u64("ratelimit-limit"));
353    let remaining = as_u64("x-ratelimit-remaining-requests")
354        .or_else(|| as_u64("x-ratelimit-remaining"))
355        .or_else(|| as_u64("ratelimit-remaining"));
356    let reset_seconds = as_f64("x-ratelimit-reset-requests")
357        .or_else(|| as_f64("x-ratelimit-reset"))
358        .or_else(|| as_f64("ratelimit-reset"));
359    (limit.is_some() || remaining.is_some() || reset_seconds.is_some())
360        .then_some(RateLimit { limit, remaining, reset_seconds })
361}
362
363/// Every header name and value, redacting anything that could be a
364/// credential by name - a bearer token echoed back, a cookie, anything with
365/// "key"/"token"/"secret" in it. `apps/jevprobe --headers` is the one caller
366/// that prints this; `Answered` always carries it so that mode needs no
367/// separate request.
368fn collect_headers(headers: &http::HeaderMap) -> Vec<(String, String)> {
369    let sensitive = |name: &str| {
370        let lower = name.to_ascii_lowercase();
371        ["authorization", "cookie", "token", "secret", "-key", "api-key"]
372            .iter()
373            .any(|word| lower.contains(word))
374    };
375    headers
376        .iter()
377        .map(|(name, value)| {
378            let name = name.as_str().to_owned();
379            let value = if sensitive(&name) {
380                "«redacted»".to_owned()
381            } else {
382                value.to_str().unwrap_or("«non-utf8»").to_owned()
383            };
384            (name, value)
385        })
386        .collect()
387}
388
389#[cfg(test)]
390mod tests {
391    use super::*;
392
393    fn model() -> ModelId {
394        ModelId::pinned("jev-1.13.0").expect("a pinned model")
395    }
396
397    fn sandbox_ledger() -> Ledger {
398        let dir = std::env::temp_dir()
399            .join(format!("jev-http-test-{}", std::process::id()))
400            .join(ledger::now().to_string());
401        Ledger::at(dir).expect("opening a sandbox ledger")
402    }
403
404    #[test]
405    fn a_missing_key_is_not_an_error() {
406        // Whatever this machine's environment happens to hold, an empty one
407        // means "no client", not "a broken client".
408        assert!(
409            Jev::new("", "test", sandbox_ledger(), model(), Endpoint::api()).is_ok(),
410            "an empty key still builds; the API refuses it"
411        );
412    }
413
414    #[test]
415    fn rate_limit_headers_are_read_when_present_and_absent_otherwise() {
416        let empty = http::HeaderMap::new();
417        assert_eq!(parse_rate_limit(&empty), None, "no header, no snapshot");
418
419        let mut headers = http::HeaderMap::new();
420        headers.insert("x-ratelimit-remaining-requests", "37".parse().expect("a header value"));
421        let rate_limit = parse_rate_limit(&headers).expect("a snapshot");
422        assert_eq!(rate_limit.remaining, Some(37));
423        assert_eq!(rate_limit.limit, None, "only what was actually sent is filled in");
424    }
425
426    #[test]
427    fn a_bearer_token_echoed_back_would_be_redacted() {
428        let mut headers = http::HeaderMap::new();
429        headers.insert("x-typesafe-request-id", "req_abc123".parse().expect("a header value"));
430        headers.insert("set-cookie", "session=deadbeef".parse().expect("a header value"));
431        let collected = collect_headers(&headers);
432        let value_of = |name: &str| {
433            collected.iter().find(|(n, _)| n == name).map(|(_, v)| v.clone()).unwrap()
434        };
435        assert_eq!(value_of("x-typesafe-request-id"), "req_abc123");
436        assert_eq!(value_of("set-cookie"), "«redacted»");
437    }
438
439    #[test]
440    fn a_huge_state_is_capped_at_the_documented_ceiling_not_left_unbounded() {
441        let huge = Json::text(&"x".repeat(1_000_000));
442        let mut questions = Questions::new();
443        questions.noul("q", jev_protocol::Noul::new(Json::text("long?"))).expect("a question");
444        assert_eq!(worst_case_tokens(&model(), &huge, &questions), 64_000.0);
445    }
446}