Skip to main content

rig_core/providers/gemini/
caching.rs

1//! Automatic explicit caching for long Gemini runs.
2//!
3//! [`Caching`](super::Caching) is a transport that wraps another. Every `generateContent`
4//! and `streamGenerateContent` request passes through it with its encoded
5//! body, and a shared [`CacheBook`] decides what the request reads from a
6//! cache and when a new cache is worth creating:
7//!
8//! ```no_run
9//! use rig_core::providers::gemini::{AutoCache, CacheBook, Gemini};
10//!
11//! # async fn example() -> Result<(), Box<dyn std::error::Error>> {
12//! let gemini = Gemini::from_env()?;
13//! let book = CacheBook::new(AutoCache::default());
14//! let model = gemini.completion("gemini-3.8-flash").caching(&book);
15//! // ... run an agent or a chat on `model` ...
16//! book.close(&gemini.cached_contents()).await;
17//! println!("{:?}", book.report());
18//! # Ok(())
19//! # }
20//! ```
21//!
22//! # What it does
23//!
24//! Gemini's implicit caching reuses an earlier request's prefix at no cost,
25//! but on gemini-3.8-flash a chat-sized request that extends the previous one
26//! almost never hits (measured: 1 hit in 138). An explicit cache
27//! (`cachedContents`) can hold the whole conversation so far: system
28//! instruction, tools, tool config, every model turn with its signatures, and
29//! function calls and responses. A request that reads it sends only what
30//! came after. Caches cannot be extended or chained, so a growing
31//! conversation needs a new cache from time to time ("rolling"), and each new
32//! cache is billed again at the input price.
33//!
34//! The `Auto` policy rolls when rolling has already paid for itself: it
35//! counts the input premium the conversation paid for tokens a cache could
36//! have held, and creates a new cache once that premium covers the new
37//! cache's creation and storage (ski rental). So:
38//!
39//! - a cache read once never pays, and is never made;
40//! - rolling beats one fixed prefix cache from about 20 calls, and the saving
41//!   grows with the run;
42//! - where implicit caching already serves a conversation (large contexts,
43//!   repeated documents), the book stays inline instead of paying to cache
44//!   what is already cheap. It goes inline above `implicit_ceiling` tokens
45//!   until implicit caching is seen failing.
46//!
47//! Cost savings top out below 90%, because a cached read costs 10% of the
48//! input price. How close a run gets depends on its shape: a large stable
49//! prefix with little new content per call approaches 99% cached (a batch of
50//! short questions over a 5.5k-token prefix measured 99.6%); a chat on a small
51//! preamble cannot, because each call's new turn is never cached, and on
52//! Gemini 3 each earlier thought a request replays through its signature is
53//! billed again as input that no explicit cache holds (see
54//! [`ThoughtReplay`](super::completion::ThoughtReplay)).
55//!
56//! # Lifecycle
57//!
58//! - The book tracks each cache's expiry and never sends one it knows has
59//!   expired. A new cache's TTL follows the gaps between calls (at least
60//!   twice the longest gap seen, never less than [`AutoCache::ttl`], never more
61//!   than 81 minutes, past which keeping a cache costs more than it saves at
62//!   the default prices), and a cache read close to expiry is extended.
63//! - A replaced cache is deleted when the roll replaces it, and a
64//!   conversation cache nobody has read for a while is deleted too.
65//! - A cache that answers 403 (deleted elsewhere, or expired early) is
66//!   forgotten and the request is sent again inline, once.
67//! - [`CacheBook::close`] deletes everything still live.
68//!
69//! # Resume
70//!
71//! A checkpoint stores [`CacheBook::leases`] beside the run. A new book in a
72//! new process calls [`CacheBook::restore`] with them, then
73//! [`CacheBook::prove`], which keeps a lease only if Google still has it
74//! under its name. The resumed run then reads the surviving cache instead of
75//! creating a new one.
76//!
77//! # Scope
78//!
79//! This covers `generateContent` and `streamGenerateContent` on the Gemini
80//! Developer API. The Interactions API has no explicit caching. Vertex AI
81//! and `rig-gemini-grpc` are not covered, nor are hosted-tool parts, which
82//! this crate's Gemini decoder does not keep in history.
83
84use std::collections::{BTreeMap, HashMap, HashSet};
85use std::sync::{Arc, Mutex, PoisonError};
86use std::time::Duration;
87
88use serde::{Deserialize, Serialize};
89use serde_json::value::RawValue;
90use sha2::{Digest, Sha256};
91
92use crate::completion::CacheCost;
93
94/// Unix seconds.
95pub type Clock = Arc<dyn Fn() -> u64 + Send + Sync>;
96
97/// Seconds since the Unix epoch, from the system clock where the target has
98/// one. On `wasm32-unknown-unknown` there is none: pass a clock with
99/// [`CacheBook::with_clock`] there, or every cache looks forever young and
100/// expiry is found by Gemini's 403 instead.
101fn system_now() -> u64 {
102    #[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
103    {
104        std::time::SystemTime::now()
105            .duration_since(std::time::UNIX_EPOCH)
106            .map_or(0, |elapsed| elapsed.as_secs())
107    }
108    #[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
109    {
110        0
111    }
112}
113
114/// The longest TTL the book gives a cache: past about 81 minutes between
115/// reads, storage costs more than a cached read saves at the default prices.
116const MAX_TTL_SECS: u64 = 81 * 60;
117
118/// A cache is used only with more than this many seconds left.
119const EXPIRY_MARGIN_SECS: u64 = 5;
120
121/// A conversation cache unread for this long, or three times its line's
122/// longest gap if longer, belongs to a conversation that moved on.
123const MIN_IDLE_SECS: u64 = 60;
124
125/// The `Auto` policy's parameters. The price ratios default to
126/// gemini-3.8-flash standard prices relative to its input price (cached read
127/// $0.075 and storage $0.50 per 1M tokens per hour, against $0.75 input);
128/// [`AutoCache::for_model`] reads the cached-read ratio of another model
129/// from the built-in catalog. The catalog holds no storage price.
130#[derive(Clone, Copy, Debug, PartialEq)]
131pub struct AutoCache {
132    /// The shortest TTL a new cache gets.
133    pub ttl: Duration,
134    /// Smallest request implicit caching serves (4,096 tokens on Gemini 3).
135    pub implicit_floor: u64,
136    /// Above this many cacheable tokens the book sends requests inline and
137    /// lets implicit caching serve them, until implicit caching is seen
138    /// failing.
139    pub implicit_ceiling: u64,
140    /// Gemini's minimum cache size.
141    pub min_tokens: u64,
142    /// The fewest tokens a new cache must add over the one it replaces.
143    pub min_gain: u64,
144    /// Cached-read price over input price.
145    pub cached_ratio: f64,
146    /// Storage price per hour over input price.
147    pub storage_ratio_per_hour: f64,
148}
149
150impl Default for AutoCache {
151    fn default() -> Self {
152        Self {
153            ttl: Duration::from_secs(60 * 60),
154            implicit_floor: 4_096,
155            implicit_ceiling: 16_000,
156            min_tokens: 1_024,
157            min_gain: 256,
158            cached_ratio: 0.1,
159            storage_ratio_per_hour: 0.5 / 0.75,
160        }
161    }
162}
163
164impl AutoCache {
165    /// The default policy with the cached-read ratio of the Gemini API's
166    /// `model` from the built-in catalog's pricing. A model the catalog
167    /// does not price keeps the default ratio.
168    pub fn for_model(model: &str) -> Self {
169        Self::priced(
170            crate::catalog::lookup(super::PROVIDER_NAME, model)
171                .and_then(|spec| spec.pricing.as_ref()),
172        )
173    }
174
175    /// The default policy with the cached-read ratio of `pricing`. Without
176    /// a cached-read price, with a negative one, or with no positive input
177    /// price to divide by, the default ratio stays.
178    fn priced(pricing: Option<&crate::catalog::Pricing>) -> Self {
179        let ratio = pricing
180            .and_then(|pricing| Some((pricing.cache_read?, pricing.input)))
181            .filter(|(read, input)| *read >= 0.0 && *input > 0.0)
182            .map(|(read, input)| read / input);
183        let default = Self::default();
184        Self {
185            cached_ratio: ratio.unwrap_or(default.cached_ratio),
186            ..default
187        }
188    }
189}
190
191/// One cache the book holds.
192#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
193pub struct Lease {
194    /// `cachedContents/<id>`.
195    pub name: String,
196    /// Hex SHA-256 of the prefix the cache holds. Its first 40 characters
197    /// end the cache's display name, which is how [`CacheBook::prove`]
198    /// recognizes it.
199    pub digest: String,
200    /// How many contents the cache holds.
201    pub covers: usize,
202    /// Its size, as Gemini counted it when it was created.
203    pub tokens: u64,
204    /// Unix seconds.
205    pub expires_at: u64,
206    /// The model the cache belongs to.
207    pub model: String,
208}
209
210/// Something that happened to a cache.
211#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
212pub enum CacheEvent {
213    /// A cache was created.
214    Created {
215        /// `cachedContents/<id>`.
216        name: String,
217        /// Its size, as Gemini counted it.
218        tokens: u64,
219        /// The book's estimate of its size before creating it.
220        estimated: u64,
221        /// How many contents it holds.
222        covers: usize,
223        /// Its TTL.
224        ttl_secs: u64,
225        /// Unix seconds.
226        at: u64,
227    },
228    /// The book deleted a cache it no longer needs.
229    Retired {
230        /// `cachedContents/<id>`.
231        name: String,
232        /// Unix seconds.
233        at: u64,
234    },
235    /// A cache answered 403 when a request named it, and was forgotten.
236    Lost {
237        /// `cachedContents/<id>`.
238        name: String,
239        /// Unix seconds.
240        at: u64,
241    },
242    /// A cache's expiry was pushed back.
243    Extended {
244        /// `cachedContents/<id>`.
245        name: String,
246        /// The new expiry, Unix seconds.
247        expires_at: u64,
248        /// Unix seconds.
249        at: u64,
250    },
251    /// Gemini refused to create a cache; the request went out without it.
252    CreateFailed {
253        /// The HTTP status, when there was one.
254        status: Option<u16>,
255        /// Gemini's message.
256        message: String,
257        /// Unix seconds.
258        at: u64,
259    },
260}
261
262/// One cache the book created.
263#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
264pub struct CreatedCache {
265    /// `cachedContents/<id>`.
266    pub name: String,
267    /// Its size.
268    pub tokens: u64,
269    /// Unix seconds.
270    pub at: u64,
271}
272
273/// What the book's caches did and cost, for pricing a run. Creation is
274/// billed like input (`tokens` × input price); storage is `token_hours` ×
275/// the storage price.
276#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
277pub struct CacheReport {
278    /// Every cache created, in order.
279    pub created: Vec<CreatedCache>,
280    /// How many requests read each cache, by name.
281    pub reads: BTreeMap<String, u64>,
282    /// Caches the book deleted because a roll replaced them, their
283    /// conversation moved on, or [`CacheBook::close`] ran.
284    pub retired: Vec<String>,
285    /// Caches that answered 403 and were forgotten.
286    pub lost: Vec<String>,
287    /// Σ tokens × hours each cache lived, until it was deleted, lost or
288    /// expired, or until now for a live one.
289    pub token_hours: f64,
290    /// The caches still live.
291    pub live: Vec<Lease>,
292}
293
294/// The book's caches as billed beside the calls: every created token as a
295/// cache write (priced at the input price) and `token_hours` of storage. Add
296/// it to the calls' [`CacheCost::from_usage`] to price a run.
297impl From<&CacheReport> for CacheCost {
298    fn from(report: &CacheReport) -> Self {
299        Self {
300            cache_writes: report.created.iter().map(|created| created.tokens).sum(),
301            storage_token_hours: report.token_hours,
302            ..Self::default()
303        }
304    }
305}
306
307#[derive(Clone, Debug, Default)]
308struct Line {
309    /// Premium tokens: cacheable input paid at full price, weighted by
310    /// `1 - cached_ratio`, since the line last rolled.
311    premium: f64,
312    /// EWMA of implicit coverage on inline calls big enough for it.
313    implicit: Option<f64>,
314    rolled_at: u64,
315    last_call: Option<u64>,
316    longest_gap: u64,
317    /// Coverable size at the last refused creation: no retry until the line
318    /// grows past it.
319    failed_at: Option<u64>,
320}
321
322#[derive(Clone, Debug)]
323struct Life {
324    tokens: u64,
325    created: u64,
326    expires_at: u64,
327    ended: Option<u64>,
328    last_read: u64,
329    line_gap: u64,
330}
331
332#[derive(Default)]
333pub(crate) struct Book {
334    pub(crate) leases: HashMap<String, Lease>,
335    lives: BTreeMap<String, Life>,
336    /// For each prefix, the conversations seen on it, by their first content.
337    lineages: HashMap<String, HashSet<String>>,
338    lines: HashMap<String, Line>,
339    events: Vec<CacheEvent>,
340    report: CacheReport,
341    /// Gemini's tokens over this book's estimate, learned from each cache it
342    /// creates.
343    calibration: Option<f64>,
344}
345
346/// The shared, content-addressed record of the caches a set of models
347/// reads. Clone it into every model that should share caches: sub-agents on
348/// one prefix, a resumed run and its successor.
349#[derive(Clone)]
350pub struct CacheBook {
351    inner: Arc<Mutex<Book>>,
352    /// Serializes cache creation, so a prefix is created once.
353    pub(crate) create: Arc<futures::lock::Mutex<()>>,
354    policy: AutoCache,
355    clock: Clock,
356    pub(crate) display_prefix: Arc<str>,
357}
358
359impl std::fmt::Debug for CacheBook {
360    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
361        f.debug_struct("CacheBook")
362            .field("policy", &self.policy)
363            .field("display_prefix", &self.display_prefix)
364            .finish_non_exhaustive()
365    }
366}
367
368impl CacheBook {
369    /// A book with `policy`, on the system clock.
370    pub fn new(policy: AutoCache) -> Self {
371        Self {
372            inner: Arc::default(),
373            create: Arc::default(),
374            policy,
375            clock: Arc::new(system_now),
376            display_prefix: Arc::from("rig-cache-"),
377        }
378    }
379
380    /// The same book reading time from `clock` (Unix seconds).
381    pub fn with_clock(mut self, clock: impl Fn() -> u64 + Send + Sync + 'static) -> Self {
382        self.clock = Arc::new(clock);
383        self
384    }
385
386    /// The same book naming its caches `<prefix><digest>`, so they can be
387    /// listed and swept. Defaults to `rig-cache-`.
388    pub fn with_display_prefix(mut self, prefix: &str) -> Self {
389        self.display_prefix = Arc::from(prefix);
390        self
391    }
392
393    /// The book's policy.
394    pub fn policy(&self) -> AutoCache {
395        self.policy
396    }
397
398    pub(crate) fn book(&self) -> std::sync::MutexGuard<'_, Book> {
399        self.inner.lock().unwrap_or_else(PoisonError::into_inner)
400    }
401
402    pub(crate) fn now(&self) -> u64 {
403        (self.clock)()
404    }
405
406    /// Everything that happened to the book's caches, in order.
407    pub fn events(&self) -> Vec<CacheEvent> {
408        self.book().events.clone()
409    }
410
411    /// The caches the book holds, for a checkpoint.
412    pub fn leases(&self) -> Vec<Lease> {
413        let mut leases: Vec<Lease> = self.book().leases.values().cloned().collect();
414        leases.sort_by(|a, b| a.name.cmp(&b.name));
415        leases
416    }
417
418    /// Put checkpointed leases back. Call [`Self::prove`] next.
419    pub fn restore(&self, leases: Vec<Lease>) {
420        let now = self.now();
421        let mut book = self.book();
422        for lease in leases {
423            book.lives.entry(lease.name.clone()).or_insert(Life {
424                tokens: lease.tokens,
425                created: now,
426                expires_at: lease.expires_at,
427                ended: None,
428                last_read: now,
429                line_gap: 0,
430            });
431            book.leases.insert(lease.digest.clone(), lease);
432        }
433    }
434
435    /// What the caches did and cost so far.
436    pub fn report(&self) -> CacheReport {
437        let now = self.now();
438        let book = self.book();
439        let mut report = book.report.clone();
440        report.token_hours = book
441            .lives
442            .values()
443            .map(|life| {
444                let end = life.ended.unwrap_or(now).min(life.expires_at);
445                life.tokens as f64 * end.saturating_sub(life.created) as f64 / 3600.0
446            })
447            .sum();
448        report.live = book.leases.values().cloned().collect();
449        report.live.sort_by(|a, b| a.name.cmp(&b.name));
450        report
451    }
452
453    pub(crate) fn retired(&self, lease: &Lease) {
454        let now = self.now();
455        let mut book = self.book();
456        book.leases.remove(&lease.digest);
457        end_life(&mut book, &lease.name, now);
458        book.events.push(CacheEvent::Retired {
459            name: lease.name.clone(),
460            at: now,
461        });
462        book.report.retired.push(lease.name.clone());
463        tracing::info!(target: "gemini.cache.retired", name = %lease.name);
464    }
465
466    pub(crate) fn lost(&self, lease: &Lease) {
467        let now = self.now();
468        let mut book = self.book();
469        book.leases.remove(&lease.digest);
470        end_life(&mut book, &lease.name, now);
471        book.events.push(CacheEvent::Lost {
472            name: lease.name.clone(),
473            at: now,
474        });
475        book.report.lost.push(lease.name.clone());
476        tracing::info!(target: "gemini.cache.lost", name = %lease.name);
477    }
478}
479
480pub(crate) fn end_life(book: &mut Book, name: &str, now: u64) {
481    if let Some(life) = book.lives.get_mut(name)
482        && life.ended.is_none()
483    {
484        life.ended = Some(now);
485    }
486}
487
488pub(crate) fn short_digest(digest: &str) -> &str {
489    digest.get(..40).unwrap_or(digest)
490}
491
492// ---------------------------------------------------------------------------
493// Request bodies: parsing, digests, stripping and cache bodies.
494
495const PREFIX_KEYS: [&str; 3] = ["systemInstruction", "tools", "toolConfig"];
496
497/// A `generateContent` body, its contents kept as their exact bytes.
498pub(crate) struct Parsed {
499    /// Every other top-level field, in order, as exact bytes.
500    pub(crate) rest: Vec<(String, Box<RawValue>)>,
501    /// `systemInstruction`, `tools` and `toolConfig`, when present and not null.
502    pub(crate) prefix: [Option<Box<RawValue>>; 3],
503    pub(crate) contents: Vec<Box<RawValue>>,
504    pub(crate) has_cached_content: bool,
505}
506
507pub(crate) fn parse(bytes: &[u8]) -> Option<Parsed> {
508    let fields: serde_json::Map<String, serde_json::Value> = serde_json::from_slice(bytes).ok()?;
509    // Parse again as raw values so every part keeps its bytes; the first
510    // pass only fixes the key order, which `RawValue` maps do not keep.
511    let raw: HashMap<String, Box<RawValue>> = serde_json::from_slice(bytes).ok()?;
512    let mut parsed = Parsed {
513        rest: Vec::new(),
514        prefix: [None, None, None],
515        contents: Vec::new(),
516        has_cached_content: false,
517    };
518    for key in fields.keys() {
519        let value = raw.get(key)?.clone();
520        if key == "contents" {
521            parsed.contents = serde_json::from_str(value.get()).ok()?;
522        } else if let Some(slot) = PREFIX_KEYS.iter().position(|prefix| prefix == key) {
523            if value.get() != "null"
524                && let Some(field) = parsed.prefix.get_mut(slot)
525            {
526                *field = Some(value);
527            }
528        } else {
529            if key == "cachedContent" && value.get() != "null" {
530                parsed.has_cached_content = true;
531            }
532            parsed.rest.push((key.clone(), value));
533        }
534    }
535    Some(parsed)
536}
537
538fn hex(bytes: &[u8]) -> String {
539    use std::fmt::Write;
540    bytes.iter().fold(String::new(), |mut out, byte| {
541        let _ = write!(out, "{byte:02x}");
542        out
543    })
544}
545
546/// `fragment` with its keys sorted: a history replays a provider part with
547/// the key order it arrived in, which a recording does not keep, so the
548/// digests read the content, not its spelling.
549fn canonical(fragment: &str) -> std::borrow::Cow<'_, str> {
550    serde_json::from_str::<serde_json::Value>(fragment)
551        .map_or(std::borrow::Cow::Borrowed(fragment), |value| {
552            std::borrow::Cow::Owned(crate::json_utils::to_canonical_string(&value))
553        })
554}
555
556/// `d[k]`: the digest of the model, the prefix and `contents[..k]`.
557pub(crate) fn digests(model: &str, parsed: &Parsed) -> Vec<String> {
558    let mut hasher = Sha256::new();
559    hasher.update(model.as_bytes());
560    for field in &parsed.prefix {
561        hasher.update([0u8]);
562        if let Some(raw) = field {
563            hasher.update(canonical(raw.get()).as_bytes());
564        }
565    }
566    let mut state = hasher.finalize().to_vec();
567    let mut out = vec![hex(&state)];
568    for content in &parsed.contents {
569        let mut hasher = Sha256::new();
570        hasher.update(&state);
571        hasher.update(canonical(content.get()).as_bytes());
572        state = hasher.finalize().to_vec();
573        out.push(hex(&state));
574    }
575    out
576}
577
578/// A rough token count of a JSON fragment: its bytes over four, without
579/// thought signatures, which Gemini expands into restored thoughts that no
580/// explicit cache holds.
581fn estimate(json: &str) -> u64 {
582    let Ok(mut value) = serde_json::from_str::<serde_json::Value>(json) else {
583        return (json.len() / 4) as u64;
584    };
585    fn scrub(value: &mut serde_json::Value) {
586        match value {
587            serde_json::Value::Object(map) => {
588                map.shift_remove("thoughtSignature");
589                map.values_mut().for_each(scrub);
590            }
591            serde_json::Value::Array(items) => items.iter_mut().for_each(scrub),
592            _ => {}
593        }
594    }
595    scrub(&mut value);
596    (value.to_string().len() / 4) as u64
597}
598
599/// Whether a content is the user's text: a new user turn, not a tool result.
600pub(crate) fn is_user_text(raw: &RawValue) -> bool {
601    #[derive(Deserialize)]
602    struct Content {
603        role: Option<String>,
604        #[serde(default)]
605        parts: Vec<serde_json::Map<String, serde_json::Value>>,
606    }
607    serde_json::from_str::<Content>(raw.get()).is_ok_and(|content| {
608        content.role.as_deref() == Some("user")
609            && content.parts.iter().any(|part| part.contains_key("text"))
610            && !content
611                .parts
612                .iter()
613                .any(|part| part.contains_key("functionResponse"))
614    })
615}
616
617/// The body that reads `lease` instead of what it holds.
618pub(crate) fn stripped(parsed: &Parsed, lease: &Lease) -> Option<Vec<u8>> {
619    let mut out: Vec<(String, &RawValue)> = Vec::new();
620    let name = serde_json::value::to_raw_value(&lease.name).ok()?;
621    let contents = serde_json::value::to_raw_value(&parsed.contents.get(lease.covers..)?).ok()?;
622    out.push(("cachedContent".to_owned(), &name));
623    out.push(("contents".to_owned(), &contents));
624    for (key, value) in &parsed.rest {
625        if key != "cachedContent" {
626            out.push((key.clone(), value));
627        }
628    }
629    let mut bytes = b"{".to_vec();
630    for (index, (key, value)) in out.iter().enumerate() {
631        if index > 0 {
632            bytes.push(b',');
633        }
634        bytes.extend(serde_json::to_vec(key).ok()?);
635        bytes.push(b':');
636        bytes.extend(value.get().as_bytes());
637    }
638    bytes.push(b'}');
639    Some(bytes)
640}
641
642/// The `cachedContents` create body for the first `covers` contents, from
643/// the request's own bytes.
644pub(crate) fn cache_body(
645    model: &str,
646    parsed: &Parsed,
647    covers: usize,
648    display_name: &str,
649    ttl_secs: u64,
650) -> Option<Vec<u8>> {
651    let mut bytes = b"{".to_vec();
652    let mut field = |key: &str, value: &str| {
653        if bytes.len() > 1 {
654            bytes.push(b',');
655        }
656        bytes.extend(format!("\"{key}\":").as_bytes());
657        bytes.extend(value.as_bytes());
658    };
659    field(
660        "model",
661        &serde_json::to_string(&format!("models/{model}")).ok()?,
662    );
663    field("displayName", &serde_json::to_string(display_name).ok()?);
664    field("ttl", &format!("\"{ttl_secs}s\""));
665    for (key, value) in PREFIX_KEYS.iter().zip(&parsed.prefix) {
666        if let Some(value) = value {
667            field(key, value.get());
668        }
669    }
670    if covers > 0 {
671        let contents = serde_json::value::to_raw_value(&parsed.contents.get(..covers)?).ok()?;
672        field("contents", contents.get());
673    }
674    bytes.push(b'}');
675    Some(bytes)
676}
677
678// ---------------------------------------------------------------------------
679// The plan.
680
681/// What one request reads, creates, retires and extends.
682pub(crate) struct Plan {
683    pub(crate) read: Option<Lease>,
684    pub(crate) create: Option<Create>,
685    pub(crate) retire: Vec<Lease>,
686    pub(crate) extend: Option<(Lease, u64)>,
687    pub(crate) line: String,
688    pub(crate) coverable: u64,
689}
690
691pub(crate) struct Create {
692    pub(crate) covers: usize,
693    pub(crate) digest: String,
694    pub(crate) ttl_secs: u64,
695    /// Uncalibrated bytes/4 estimate of what the cache holds.
696    estimate: u64,
697    /// The calibrated estimate the plan used.
698    expected: u64,
699    pub(crate) replaces: Option<Lease>,
700}
701
702impl CacheBook {
703    pub(crate) fn plan(
704        &self,
705        model: &str,
706        d: &[String],
707        parsed: &Parsed,
708        roll_allowed: bool,
709    ) -> Plan {
710        let p = self.policy;
711        let now = self.now();
712        let n = parsed.contents.len();
713        // `d` holds one digest per content boundary, `0..=n`.
714        let at = |k: usize| d.get(k).cloned().unwrap_or_default();
715        let mut book = self.book();
716        let calibration = book.calibration.unwrap_or(1.0);
717        let prefix_estimate: u64 = parsed
718            .prefix
719            .iter()
720            .flatten()
721            .map(|raw| estimate(raw.get()))
722            .sum();
723        let content_estimates: Vec<u64> = parsed
724            .contents
725            .iter()
726            .map(|raw| estimate(raw.get()))
727            .collect();
728
729        // Who shares this prefix: a second conversation on it (sub-agents,
730        // concurrent users, or the conversation a compaction started).
731        if n >= 1 {
732            book.lineages.entry(at(0)).or_default().insert(at(1));
733        }
734        let shared = book
735            .lineages
736            .get(&at(0))
737            .is_some_and(|seen| seen.len() >= 2);
738
739        // Retire conversation caches nobody reads any more.
740        let idle: Vec<Lease> = book
741            .leases
742            .values()
743            .filter(|lease| lease.covers > 0 && lease.model == model)
744            .filter(|lease| {
745                book.lives.get(&lease.name).is_some_and(|life| {
746                    now.saturating_sub(life.last_read)
747                        > MIN_IDLE_SECS.max(life.line_gap.saturating_mul(3))
748                })
749            })
750            .cloned()
751            .collect();
752
753        // The longest live cache this request starts with; never the newest content.
754        let read = (0..n.max(1))
755            .rev()
756            .find_map(|k| {
757                book.leases
758                    .get(&at(k))
759                    .filter(|lease| lease.model == model)
760                    .filter(|lease| lease.expires_at > now + EXPIRY_MARGIN_SECS)
761                    .cloned()
762            })
763            .filter(|lease| !idle.iter().any(|gone| gone.name == lease.name));
764        let covered = read.as_ref().map_or(0, |lease| lease.covers);
765        let tail: u64 = content_estimates
766            .get(covered..n.saturating_sub(1).max(covered))
767            .unwrap_or_default()
768            .iter()
769            .sum();
770        let coverable = match &read {
771            Some(lease) => lease.tokens + (tail as f64 * calibration) as u64,
772            None => ((prefix_estimate + tail) as f64 * calibration) as u64,
773        };
774
775        // The conversation's line: its cache, or its first content when inline.
776        let line_key = read
777            .as_ref()
778            .filter(|lease| lease.covers > 0)
779            .map_or_else(|| at(1.min(n)), |lease| lease.digest.clone());
780        let prefix_leased = book.leases.contains_key(&at(0));
781        let line = book.lines.entry(line_key.clone()).or_insert_with(|| Line {
782            rolled_at: now,
783            ..Line::default()
784        });
785        if let Some(last) = line.last_call {
786            line.longest_gap = line.longest_gap.max(now.saturating_sub(last));
787        }
788        line.last_call = Some(now);
789        let longest_gap = line.longest_gap;
790        let ttl_secs = p
791            .ttl
792            .as_secs()
793            .max(longest_gap.saturating_mul(2).min(MAX_TTL_SECS))
794            .max(1);
795
796        let mut plan = Plan {
797            read: read.clone(),
798            create: None,
799            retire: idle,
800            extend: None,
801            line: line_key,
802            coverable,
803        };
804
805        // Past the ceiling, implicit caching serves: go inline until it fails.
806        let implicit_serving = line.implicit.is_none_or(|ratio| ratio >= 0.5);
807        if coverable >= p.implicit_ceiling && implicit_serving {
808            plan.read = None;
809            return plan;
810        }
811
812        let failed_below = line.failed_at.is_some_and(|at| coverable < at + p.min_gain);
813        if read.is_none()
814            && shared
815            && !prefix_leased
816            && (prefix_estimate as f64 * calibration) as u64 >= p.min_tokens
817            && !failed_below
818        {
819            plan.create = Some(Create {
820                covers: 0,
821                digest: at(0),
822                ttl_secs,
823                estimate: prefix_estimate,
824                expected: (prefix_estimate as f64 * calibration) as u64,
825                replaces: None,
826            });
827        } else if n >= 2 && roll_allowed && !failed_below {
828            let cached = read.as_ref().map_or(0, |lease| lease.tokens);
829            let gain = coverable.saturating_sub(cached);
830            let held_hours = now.saturating_sub(line.rolled_at).max(60) as f64 / 3600.0;
831            let cost = coverable as f64 * (1.0 + p.storage_ratio_per_hour * held_hours);
832            if gain >= p.min_gain && coverable >= p.min_tokens && line.premium >= cost {
833                plan.create = Some(Create {
834                    covers: n - 1,
835                    digest: at(n - 1),
836                    ttl_secs,
837                    estimate: prefix_estimate
838                        + content_estimates
839                            .get(..n - 1)
840                            .unwrap_or_default()
841                            .iter()
842                            .sum::<u64>(),
843                    expected: coverable,
844                    replaces: read.clone().filter(|lease| lease.covers > 0),
845                });
846            }
847        }
848
849        // Extend a cache that would expire before its line's next call.
850        if plan.create.is_none()
851            && let Some(lease) = &read
852            && longest_gap > 0
853            && lease.expires_at.saturating_sub(now) < longest_gap.saturating_mul(3) / 2
854        {
855            plan.extend = Some((lease.clone(), ttl_secs));
856        }
857        plan
858    }
859
860    /// Account for one reply's usage on `line`.
861    pub(crate) fn observe(&self, line: &str, read: Option<&str>, coverable: u64, cached: u64) {
862        let p = self.policy;
863        let mut book = self.book();
864        if let Some(name) = read {
865            *book.report.reads.entry(name.to_owned()).or_default() += 1;
866        }
867        let entry = book.lines.entry(line.to_owned()).or_default();
868        let mut missed = coverable.saturating_sub(cached) as f64;
869        if read.is_none() && coverable >= p.implicit_floor {
870            let ratio = cached as f64 / coverable.max(1) as f64;
871            let implicit = entry
872                .implicit
873                .map_or(ratio, |ewma| 0.5 * ewma + 0.5 * ratio);
874            entry.implicit = Some(implicit);
875            if implicit >= 0.5 {
876                missed = 0.0;
877            }
878        }
879        entry.premium += missed * (1.0 - p.cached_ratio);
880    }
881
882    pub(crate) fn created(
883        &self,
884        from_line: &str,
885        create: &Create,
886        name: String,
887        tokens: u64,
888        model: &str,
889    ) -> Lease {
890        let now = self.now();
891        let lease = Lease {
892            name: name.clone(),
893            digest: create.digest.clone(),
894            covers: create.covers,
895            tokens,
896            expires_at: now + create.ttl_secs,
897            model: model.to_owned(),
898        };
899        let mut book = self.book();
900        if create.estimate > 0 && tokens > 0 {
901            let ratio = tokens as f64 / create.estimate as f64;
902            book.calibration = Some(book.calibration.map_or(ratio, |c| 0.5 * c + 0.5 * ratio));
903        }
904        let previous = book.lines.get(from_line).cloned().unwrap_or_default();
905        if create.covers > 0 {
906            book.lines.insert(
907                lease.digest.clone(),
908                Line {
909                    premium: 0.0,
910                    implicit: previous.implicit,
911                    rolled_at: now,
912                    last_call: previous.last_call,
913                    longest_gap: previous.longest_gap,
914                    failed_at: None,
915                },
916            );
917        }
918        book.lives.insert(
919            name.clone(),
920            Life {
921                tokens,
922                created: now,
923                expires_at: lease.expires_at,
924                ended: None,
925                last_read: now,
926                line_gap: previous.longest_gap,
927            },
928        );
929        book.events.push(CacheEvent::Created {
930            name: name.clone(),
931            tokens,
932            estimated: create.expected,
933            covers: create.covers,
934            ttl_secs: create.ttl_secs,
935            at: now,
936        });
937        book.report.created.push(CreatedCache {
938            name,
939            tokens,
940            at: now,
941        });
942        book.leases.insert(lease.digest.clone(), lease.clone());
943        tracing::info!(target: "gemini.cache.created", name = %lease.name, tokens, covers = create.covers);
944        lease
945    }
946
947    pub(crate) fn create_failed(
948        &self,
949        line: &str,
950        coverable: u64,
951        status: Option<u16>,
952        message: String,
953    ) {
954        let now = self.now();
955        let mut book = self.book();
956        if let Some(entry) = book.lines.get_mut(line) {
957            entry.failed_at = Some(coverable);
958        }
959        book.events.push(CacheEvent::CreateFailed {
960            status,
961            message,
962            at: now,
963        });
964    }
965
966    pub(crate) fn extended(&self, lease: &Lease, ttl_secs: u64) {
967        let now = self.now();
968        let expires_at = now + ttl_secs;
969        let mut book = self.book();
970        if let Some(held) = book.leases.get_mut(&lease.digest) {
971            held.expires_at = expires_at;
972        }
973        if let Some(life) = book.lives.get_mut(&lease.name) {
974            life.expires_at = expires_at;
975        }
976        book.events.push(CacheEvent::Extended {
977            name: lease.name.clone(),
978            expires_at,
979            at: now,
980        });
981        tracing::info!(target: "gemini.cache.extended", name = %lease.name, expires_at);
982    }
983
984    pub(crate) fn touched(&self, lease: &Lease) {
985        let now = self.now();
986        let mut book = self.book();
987        let gap = book
988            .lines
989            .get(&lease.digest)
990            .map_or(0, |line| line.longest_gap);
991        if let Some(life) = book.lives.get_mut(&lease.name) {
992            life.last_read = now;
993            life.line_gap = life.line_gap.max(gap);
994        }
995    }
996}
997
998#[cfg(test)]
999mod tests;