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/// set them for other models or tiers. Nothing here reads a model name.
129#[derive(Clone, Copy, Debug, PartialEq)]
130pub struct AutoCache {
131    /// The shortest TTL a new cache gets.
132    pub ttl: Duration,
133    /// Smallest request implicit caching serves (4,096 tokens on Gemini 3).
134    pub implicit_floor: u64,
135    /// Above this many cacheable tokens the book sends requests inline and
136    /// lets implicit caching serve them, until implicit caching is seen
137    /// failing.
138    pub implicit_ceiling: u64,
139    /// Gemini's minimum cache size.
140    pub min_tokens: u64,
141    /// The fewest tokens a new cache must add over the one it replaces.
142    pub min_gain: u64,
143    /// Cached-read price over input price.
144    pub cached_ratio: f64,
145    /// Storage price per hour over input price.
146    pub storage_ratio_per_hour: f64,
147}
148
149impl Default for AutoCache {
150    fn default() -> Self {
151        Self {
152            ttl: Duration::from_secs(60 * 60),
153            implicit_floor: 4_096,
154            implicit_ceiling: 16_000,
155            min_tokens: 1_024,
156            min_gain: 256,
157            cached_ratio: 0.1,
158            storage_ratio_per_hour: 0.5 / 0.75,
159        }
160    }
161}
162
163/// One cache the book holds.
164#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
165pub struct Lease {
166    /// `cachedContents/<id>`.
167    pub name: String,
168    /// Hex SHA-256 of the prefix the cache holds. Its first 40 characters
169    /// end the cache's display name, which is how [`CacheBook::prove`]
170    /// recognizes it.
171    pub digest: String,
172    /// How many contents the cache holds.
173    pub covers: usize,
174    /// Its size, as Gemini counted it when it was created.
175    pub tokens: u64,
176    /// Unix seconds.
177    pub expires_at: u64,
178    /// The model the cache belongs to.
179    pub model: String,
180}
181
182/// Something that happened to a cache.
183#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
184pub enum CacheEvent {
185    /// A cache was created.
186    Created {
187        /// `cachedContents/<id>`.
188        name: String,
189        /// Its size, as Gemini counted it.
190        tokens: u64,
191        /// The book's estimate of its size before creating it.
192        estimated: u64,
193        /// How many contents it holds.
194        covers: usize,
195        /// Its TTL.
196        ttl_secs: u64,
197        /// Unix seconds.
198        at: u64,
199    },
200    /// The book deleted a cache it no longer needs.
201    Retired {
202        /// `cachedContents/<id>`.
203        name: String,
204        /// Unix seconds.
205        at: u64,
206    },
207    /// A cache answered 403 when a request named it, and was forgotten.
208    Lost {
209        /// `cachedContents/<id>`.
210        name: String,
211        /// Unix seconds.
212        at: u64,
213    },
214    /// A cache's expiry was pushed back.
215    Extended {
216        /// `cachedContents/<id>`.
217        name: String,
218        /// The new expiry, Unix seconds.
219        expires_at: u64,
220        /// Unix seconds.
221        at: u64,
222    },
223    /// Gemini refused to create a cache; the request went out without it.
224    CreateFailed {
225        /// The HTTP status, when there was one.
226        status: Option<u16>,
227        /// Gemini's message.
228        message: String,
229        /// Unix seconds.
230        at: u64,
231    },
232}
233
234/// One cache the book created.
235#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
236pub struct CreatedCache {
237    /// `cachedContents/<id>`.
238    pub name: String,
239    /// Its size.
240    pub tokens: u64,
241    /// Unix seconds.
242    pub at: u64,
243}
244
245/// What the book's caches did and cost, for pricing a run. Creation is
246/// billed like input (`tokens` × input price); storage is `token_hours` ×
247/// the storage price.
248#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
249pub struct CacheReport {
250    /// Every cache created, in order.
251    pub created: Vec<CreatedCache>,
252    /// How many requests read each cache, by name.
253    pub reads: BTreeMap<String, u64>,
254    /// Caches the book deleted because a roll replaced them, their
255    /// conversation moved on, or [`CacheBook::close`] ran.
256    pub retired: Vec<String>,
257    /// Caches that answered 403 and were forgotten.
258    pub lost: Vec<String>,
259    /// Σ tokens × hours each cache lived, until it was deleted, lost or
260    /// expired, or until now for a live one.
261    pub token_hours: f64,
262    /// The caches still live.
263    pub live: Vec<Lease>,
264}
265
266/// The book's caches as billed beside the calls: every created token as a
267/// cache write (priced at the input price) and `token_hours` of storage. Add
268/// it to the calls' [`CacheCost::from_usage`] to price a run.
269impl From<&CacheReport> for CacheCost {
270    fn from(report: &CacheReport) -> Self {
271        Self {
272            cache_writes: report.created.iter().map(|created| created.tokens).sum(),
273            storage_token_hours: report.token_hours,
274            ..Self::default()
275        }
276    }
277}
278
279#[derive(Clone, Debug, Default)]
280struct Line {
281    /// Premium tokens: cacheable input paid at full price, weighted by
282    /// `1 - cached_ratio`, since the line last rolled.
283    premium: f64,
284    /// EWMA of implicit coverage on inline calls big enough for it.
285    implicit: Option<f64>,
286    rolled_at: u64,
287    last_call: Option<u64>,
288    longest_gap: u64,
289    /// Coverable size at the last refused creation: no retry until the line
290    /// grows past it.
291    failed_at: Option<u64>,
292}
293
294#[derive(Clone, Debug)]
295struct Life {
296    tokens: u64,
297    created: u64,
298    expires_at: u64,
299    ended: Option<u64>,
300    last_read: u64,
301    line_gap: u64,
302}
303
304#[derive(Default)]
305pub(crate) struct Book {
306    pub(crate) leases: HashMap<String, Lease>,
307    lives: BTreeMap<String, Life>,
308    /// For each prefix, the conversations seen on it, by their first content.
309    lineages: HashMap<String, HashSet<String>>,
310    lines: HashMap<String, Line>,
311    events: Vec<CacheEvent>,
312    report: CacheReport,
313    /// Gemini's tokens over this book's estimate, learned from each cache it
314    /// creates.
315    calibration: Option<f64>,
316}
317
318/// The shared, content-addressed record of the caches a set of models
319/// reads. Clone it into every model that should share caches: sub-agents on
320/// one prefix, a resumed run and its successor.
321#[derive(Clone)]
322pub struct CacheBook {
323    inner: Arc<Mutex<Book>>,
324    /// Serializes cache creation, so a prefix is created once.
325    pub(crate) create: Arc<futures::lock::Mutex<()>>,
326    policy: AutoCache,
327    clock: Clock,
328    pub(crate) display_prefix: Arc<str>,
329}
330
331impl std::fmt::Debug for CacheBook {
332    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
333        f.debug_struct("CacheBook")
334            .field("policy", &self.policy)
335            .field("display_prefix", &self.display_prefix)
336            .finish_non_exhaustive()
337    }
338}
339
340impl CacheBook {
341    /// A book with `policy`, on the system clock.
342    pub fn new(policy: AutoCache) -> Self {
343        Self {
344            inner: Arc::default(),
345            create: Arc::default(),
346            policy,
347            clock: Arc::new(system_now),
348            display_prefix: Arc::from("rig-cache-"),
349        }
350    }
351
352    /// The same book reading time from `clock` (Unix seconds).
353    pub fn with_clock(mut self, clock: impl Fn() -> u64 + Send + Sync + 'static) -> Self {
354        self.clock = Arc::new(clock);
355        self
356    }
357
358    /// The same book naming its caches `<prefix><digest>`, so they can be
359    /// listed and swept. Defaults to `rig-cache-`.
360    pub fn with_display_prefix(mut self, prefix: &str) -> Self {
361        self.display_prefix = Arc::from(prefix);
362        self
363    }
364
365    /// The book's policy.
366    pub fn policy(&self) -> AutoCache {
367        self.policy
368    }
369
370    pub(crate) fn book(&self) -> std::sync::MutexGuard<'_, Book> {
371        self.inner.lock().unwrap_or_else(PoisonError::into_inner)
372    }
373
374    pub(crate) fn now(&self) -> u64 {
375        (self.clock)()
376    }
377
378    /// Everything that happened to the book's caches, in order.
379    pub fn events(&self) -> Vec<CacheEvent> {
380        self.book().events.clone()
381    }
382
383    /// The caches the book holds, for a checkpoint.
384    pub fn leases(&self) -> Vec<Lease> {
385        let mut leases: Vec<Lease> = self.book().leases.values().cloned().collect();
386        leases.sort_by(|a, b| a.name.cmp(&b.name));
387        leases
388    }
389
390    /// Put checkpointed leases back. Call [`Self::prove`] next.
391    pub fn restore(&self, leases: Vec<Lease>) {
392        let now = self.now();
393        let mut book = self.book();
394        for lease in leases {
395            book.lives.entry(lease.name.clone()).or_insert(Life {
396                tokens: lease.tokens,
397                created: now,
398                expires_at: lease.expires_at,
399                ended: None,
400                last_read: now,
401                line_gap: 0,
402            });
403            book.leases.insert(lease.digest.clone(), lease);
404        }
405    }
406
407    /// What the caches did and cost so far.
408    pub fn report(&self) -> CacheReport {
409        let now = self.now();
410        let book = self.book();
411        let mut report = book.report.clone();
412        report.token_hours = book
413            .lives
414            .values()
415            .map(|life| {
416                let end = life.ended.unwrap_or(now).min(life.expires_at);
417                life.tokens as f64 * end.saturating_sub(life.created) as f64 / 3600.0
418            })
419            .sum();
420        report.live = book.leases.values().cloned().collect();
421        report.live.sort_by(|a, b| a.name.cmp(&b.name));
422        report
423    }
424
425    pub(crate) fn retired(&self, lease: &Lease) {
426        let now = self.now();
427        let mut book = self.book();
428        book.leases.remove(&lease.digest);
429        end_life(&mut book, &lease.name, now);
430        book.events.push(CacheEvent::Retired {
431            name: lease.name.clone(),
432            at: now,
433        });
434        book.report.retired.push(lease.name.clone());
435        tracing::info!(target: "gemini.cache.retired", name = %lease.name);
436    }
437
438    pub(crate) fn lost(&self, lease: &Lease) {
439        let now = self.now();
440        let mut book = self.book();
441        book.leases.remove(&lease.digest);
442        end_life(&mut book, &lease.name, now);
443        book.events.push(CacheEvent::Lost {
444            name: lease.name.clone(),
445            at: now,
446        });
447        book.report.lost.push(lease.name.clone());
448        tracing::info!(target: "gemini.cache.lost", name = %lease.name);
449    }
450}
451
452pub(crate) fn end_life(book: &mut Book, name: &str, now: u64) {
453    if let Some(life) = book.lives.get_mut(name)
454        && life.ended.is_none()
455    {
456        life.ended = Some(now);
457    }
458}
459
460pub(crate) fn short_digest(digest: &str) -> &str {
461    digest.get(..40).unwrap_or(digest)
462}
463
464// ---------------------------------------------------------------------------
465// Request bodies: parsing, digests, stripping and cache bodies.
466
467const PREFIX_KEYS: [&str; 3] = ["systemInstruction", "tools", "toolConfig"];
468
469/// A `generateContent` body, its contents kept as their exact bytes.
470pub(crate) struct Parsed {
471    /// Every other top-level field, in order, as exact bytes.
472    pub(crate) rest: Vec<(String, Box<RawValue>)>,
473    /// `systemInstruction`, `tools` and `toolConfig`, when present and not null.
474    pub(crate) prefix: [Option<Box<RawValue>>; 3],
475    pub(crate) contents: Vec<Box<RawValue>>,
476    pub(crate) has_cached_content: bool,
477}
478
479pub(crate) fn parse(bytes: &[u8]) -> Option<Parsed> {
480    let fields: serde_json::Map<String, serde_json::Value> = serde_json::from_slice(bytes).ok()?;
481    // Parse again as raw values so every part keeps its bytes; the first
482    // pass only fixes the key order, which `RawValue` maps do not keep.
483    let raw: HashMap<String, Box<RawValue>> = serde_json::from_slice(bytes).ok()?;
484    let mut parsed = Parsed {
485        rest: Vec::new(),
486        prefix: [None, None, None],
487        contents: Vec::new(),
488        has_cached_content: false,
489    };
490    for key in fields.keys() {
491        let value = raw.get(key)?.clone();
492        if key == "contents" {
493            parsed.contents = serde_json::from_str(value.get()).ok()?;
494        } else if let Some(slot) = PREFIX_KEYS.iter().position(|prefix| prefix == key) {
495            if value.get() != "null"
496                && let Some(field) = parsed.prefix.get_mut(slot)
497            {
498                *field = Some(value);
499            }
500        } else {
501            if key == "cachedContent" && value.get() != "null" {
502                parsed.has_cached_content = true;
503            }
504            parsed.rest.push((key.clone(), value));
505        }
506    }
507    Some(parsed)
508}
509
510fn hex(bytes: &[u8]) -> String {
511    use std::fmt::Write;
512    bytes.iter().fold(String::new(), |mut out, byte| {
513        let _ = write!(out, "{byte:02x}");
514        out
515    })
516}
517
518/// `d[k]`: the digest of the model, the prefix and `contents[..k]`.
519pub(crate) fn digests(model: &str, parsed: &Parsed) -> Vec<String> {
520    let mut hasher = Sha256::new();
521    hasher.update(model.as_bytes());
522    for field in &parsed.prefix {
523        hasher.update([0u8]);
524        if let Some(raw) = field {
525            hasher.update(raw.get().as_bytes());
526        }
527    }
528    let mut state = hasher.finalize().to_vec();
529    let mut out = vec![hex(&state)];
530    for content in &parsed.contents {
531        let mut hasher = Sha256::new();
532        hasher.update(&state);
533        hasher.update(content.get().as_bytes());
534        state = hasher.finalize().to_vec();
535        out.push(hex(&state));
536    }
537    out
538}
539
540/// A rough token count of a JSON fragment: its bytes over four, without
541/// thought signatures, which Gemini expands into restored thoughts that no
542/// explicit cache holds.
543fn estimate(json: &str) -> u64 {
544    let Ok(mut value) = serde_json::from_str::<serde_json::Value>(json) else {
545        return (json.len() / 4) as u64;
546    };
547    fn scrub(value: &mut serde_json::Value) {
548        match value {
549            serde_json::Value::Object(map) => {
550                map.remove("thoughtSignature");
551                map.values_mut().for_each(scrub);
552            }
553            serde_json::Value::Array(items) => items.iter_mut().for_each(scrub),
554            _ => {}
555        }
556    }
557    scrub(&mut value);
558    (value.to_string().len() / 4) as u64
559}
560
561/// Whether a content is the user's text: a new user turn, not a tool result.
562pub(crate) fn is_user_text(raw: &RawValue) -> bool {
563    #[derive(Deserialize)]
564    struct Content {
565        role: Option<String>,
566        #[serde(default)]
567        parts: Vec<serde_json::Map<String, serde_json::Value>>,
568    }
569    serde_json::from_str::<Content>(raw.get()).is_ok_and(|content| {
570        content.role.as_deref() == Some("user")
571            && content.parts.iter().any(|part| part.contains_key("text"))
572            && !content
573                .parts
574                .iter()
575                .any(|part| part.contains_key("functionResponse"))
576    })
577}
578
579/// The body that reads `lease` instead of what it holds.
580pub(crate) fn stripped(parsed: &Parsed, lease: &Lease) -> Option<Vec<u8>> {
581    let mut out: Vec<(String, &RawValue)> = Vec::new();
582    let name = serde_json::value::to_raw_value(&lease.name).ok()?;
583    let contents = serde_json::value::to_raw_value(&parsed.contents.get(lease.covers..)?).ok()?;
584    out.push(("cachedContent".to_owned(), &name));
585    out.push(("contents".to_owned(), &contents));
586    for (key, value) in &parsed.rest {
587        if key != "cachedContent" {
588            out.push((key.clone(), value));
589        }
590    }
591    let mut bytes = b"{".to_vec();
592    for (index, (key, value)) in out.iter().enumerate() {
593        if index > 0 {
594            bytes.push(b',');
595        }
596        bytes.extend(serde_json::to_vec(key).ok()?);
597        bytes.push(b':');
598        bytes.extend(value.get().as_bytes());
599    }
600    bytes.push(b'}');
601    Some(bytes)
602}
603
604/// The `cachedContents` create body for the first `covers` contents, from
605/// the request's own bytes.
606pub(crate) fn cache_body(
607    model: &str,
608    parsed: &Parsed,
609    covers: usize,
610    display_name: &str,
611    ttl_secs: u64,
612) -> Option<Vec<u8>> {
613    let mut bytes = b"{".to_vec();
614    let mut field = |key: &str, value: &str| {
615        if bytes.len() > 1 {
616            bytes.push(b',');
617        }
618        bytes.extend(format!("\"{key}\":").as_bytes());
619        bytes.extend(value.as_bytes());
620    };
621    field(
622        "model",
623        &serde_json::to_string(&format!("models/{model}")).ok()?,
624    );
625    field("displayName", &serde_json::to_string(display_name).ok()?);
626    field("ttl", &format!("\"{ttl_secs}s\""));
627    for (key, value) in PREFIX_KEYS.iter().zip(&parsed.prefix) {
628        if let Some(value) = value {
629            field(key, value.get());
630        }
631    }
632    if covers > 0 {
633        let contents = serde_json::value::to_raw_value(&parsed.contents.get(..covers)?).ok()?;
634        field("contents", contents.get());
635    }
636    bytes.push(b'}');
637    Some(bytes)
638}
639
640// ---------------------------------------------------------------------------
641// The plan.
642
643/// What one request reads, creates, retires and extends.
644pub(crate) struct Plan {
645    pub(crate) read: Option<Lease>,
646    pub(crate) create: Option<Create>,
647    pub(crate) retire: Vec<Lease>,
648    pub(crate) extend: Option<(Lease, u64)>,
649    pub(crate) line: String,
650    pub(crate) coverable: u64,
651}
652
653pub(crate) struct Create {
654    pub(crate) covers: usize,
655    pub(crate) digest: String,
656    pub(crate) ttl_secs: u64,
657    /// Uncalibrated bytes/4 estimate of what the cache holds.
658    estimate: u64,
659    /// The calibrated estimate the plan used.
660    expected: u64,
661    pub(crate) replaces: Option<Lease>,
662}
663
664impl CacheBook {
665    pub(crate) fn plan(
666        &self,
667        model: &str,
668        d: &[String],
669        parsed: &Parsed,
670        roll_allowed: bool,
671    ) -> Plan {
672        let p = self.policy;
673        let now = self.now();
674        let n = parsed.contents.len();
675        // `d` holds one digest per content boundary, `0..=n`.
676        let at = |k: usize| d.get(k).cloned().unwrap_or_default();
677        let mut book = self.book();
678        let calibration = book.calibration.unwrap_or(1.0);
679        let prefix_estimate: u64 = parsed
680            .prefix
681            .iter()
682            .flatten()
683            .map(|raw| estimate(raw.get()))
684            .sum();
685        let content_estimates: Vec<u64> = parsed
686            .contents
687            .iter()
688            .map(|raw| estimate(raw.get()))
689            .collect();
690
691        // Who shares this prefix: a second conversation on it (sub-agents,
692        // concurrent users, or the conversation a compaction started).
693        if n >= 1 {
694            book.lineages.entry(at(0)).or_default().insert(at(1));
695        }
696        let shared = book
697            .lineages
698            .get(&at(0))
699            .is_some_and(|seen| seen.len() >= 2);
700
701        // Retire conversation caches nobody reads any more.
702        let idle: Vec<Lease> = book
703            .leases
704            .values()
705            .filter(|lease| lease.covers > 0 && lease.model == model)
706            .filter(|lease| {
707                book.lives.get(&lease.name).is_some_and(|life| {
708                    now.saturating_sub(life.last_read)
709                        > MIN_IDLE_SECS.max(life.line_gap.saturating_mul(3))
710                })
711            })
712            .cloned()
713            .collect();
714
715        // The longest live cache this request starts with; never the newest content.
716        let read = (0..n.max(1))
717            .rev()
718            .find_map(|k| {
719                book.leases
720                    .get(&at(k))
721                    .filter(|lease| lease.model == model)
722                    .filter(|lease| lease.expires_at > now + EXPIRY_MARGIN_SECS)
723                    .cloned()
724            })
725            .filter(|lease| !idle.iter().any(|gone| gone.name == lease.name));
726        let covered = read.as_ref().map_or(0, |lease| lease.covers);
727        let tail: u64 = content_estimates
728            .get(covered..n.saturating_sub(1).max(covered))
729            .unwrap_or_default()
730            .iter()
731            .sum();
732        let coverable = match &read {
733            Some(lease) => lease.tokens + (tail as f64 * calibration) as u64,
734            None => ((prefix_estimate + tail) as f64 * calibration) as u64,
735        };
736
737        // The conversation's line: its cache, or its first content when inline.
738        let line_key = read
739            .as_ref()
740            .filter(|lease| lease.covers > 0)
741            .map_or_else(|| at(1.min(n)), |lease| lease.digest.clone());
742        let prefix_leased = book.leases.contains_key(&at(0));
743        let line = book.lines.entry(line_key.clone()).or_insert_with(|| Line {
744            rolled_at: now,
745            ..Line::default()
746        });
747        if let Some(last) = line.last_call {
748            line.longest_gap = line.longest_gap.max(now.saturating_sub(last));
749        }
750        line.last_call = Some(now);
751        let longest_gap = line.longest_gap;
752        let ttl_secs = p
753            .ttl
754            .as_secs()
755            .max(longest_gap.saturating_mul(2).min(MAX_TTL_SECS))
756            .max(1);
757
758        let mut plan = Plan {
759            read: read.clone(),
760            create: None,
761            retire: idle,
762            extend: None,
763            line: line_key,
764            coverable,
765        };
766
767        // Past the ceiling, implicit caching serves: go inline until it fails.
768        let implicit_serving = line.implicit.is_none_or(|ratio| ratio >= 0.5);
769        if coverable >= p.implicit_ceiling && implicit_serving {
770            plan.read = None;
771            return plan;
772        }
773
774        let failed_below = line.failed_at.is_some_and(|at| coverable < at + p.min_gain);
775        if read.is_none()
776            && shared
777            && !prefix_leased
778            && (prefix_estimate as f64 * calibration) as u64 >= p.min_tokens
779            && !failed_below
780        {
781            plan.create = Some(Create {
782                covers: 0,
783                digest: at(0),
784                ttl_secs,
785                estimate: prefix_estimate,
786                expected: (prefix_estimate as f64 * calibration) as u64,
787                replaces: None,
788            });
789        } else if n >= 2 && roll_allowed && !failed_below {
790            let cached = read.as_ref().map_or(0, |lease| lease.tokens);
791            let gain = coverable.saturating_sub(cached);
792            let held_hours = now.saturating_sub(line.rolled_at).max(60) as f64 / 3600.0;
793            let cost = coverable as f64 * (1.0 + p.storage_ratio_per_hour * held_hours);
794            if gain >= p.min_gain && coverable >= p.min_tokens && line.premium >= cost {
795                plan.create = Some(Create {
796                    covers: n - 1,
797                    digest: at(n - 1),
798                    ttl_secs,
799                    estimate: prefix_estimate
800                        + content_estimates
801                            .get(..n - 1)
802                            .unwrap_or_default()
803                            .iter()
804                            .sum::<u64>(),
805                    expected: coverable,
806                    replaces: read.clone().filter(|lease| lease.covers > 0),
807                });
808            }
809        }
810
811        // Extend a cache that would expire before its line's next call.
812        if plan.create.is_none()
813            && let Some(lease) = &read
814            && longest_gap > 0
815            && lease.expires_at.saturating_sub(now) < longest_gap.saturating_mul(3) / 2
816        {
817            plan.extend = Some((lease.clone(), ttl_secs));
818        }
819        plan
820    }
821
822    /// Account for one reply's usage on `line`.
823    pub(crate) fn observe(&self, line: &str, read: Option<&str>, coverable: u64, cached: u64) {
824        let p = self.policy;
825        let mut book = self.book();
826        if let Some(name) = read {
827            *book.report.reads.entry(name.to_owned()).or_default() += 1;
828        }
829        let entry = book.lines.entry(line.to_owned()).or_default();
830        let mut missed = coverable.saturating_sub(cached) as f64;
831        if read.is_none() && coverable >= p.implicit_floor {
832            let ratio = cached as f64 / coverable.max(1) as f64;
833            let implicit = entry
834                .implicit
835                .map_or(ratio, |ewma| 0.5 * ewma + 0.5 * ratio);
836            entry.implicit = Some(implicit);
837            if implicit >= 0.5 {
838                missed = 0.0;
839            }
840        }
841        entry.premium += missed * (1.0 - p.cached_ratio);
842    }
843
844    pub(crate) fn created(
845        &self,
846        from_line: &str,
847        create: &Create,
848        name: String,
849        tokens: u64,
850        model: &str,
851    ) -> Lease {
852        let now = self.now();
853        let lease = Lease {
854            name: name.clone(),
855            digest: create.digest.clone(),
856            covers: create.covers,
857            tokens,
858            expires_at: now + create.ttl_secs,
859            model: model.to_owned(),
860        };
861        let mut book = self.book();
862        if create.estimate > 0 && tokens > 0 {
863            let ratio = tokens as f64 / create.estimate as f64;
864            book.calibration = Some(book.calibration.map_or(ratio, |c| 0.5 * c + 0.5 * ratio));
865        }
866        let previous = book.lines.get(from_line).cloned().unwrap_or_default();
867        if create.covers > 0 {
868            book.lines.insert(
869                lease.digest.clone(),
870                Line {
871                    premium: 0.0,
872                    implicit: previous.implicit,
873                    rolled_at: now,
874                    last_call: previous.last_call,
875                    longest_gap: previous.longest_gap,
876                    failed_at: None,
877                },
878            );
879        }
880        book.lives.insert(
881            name.clone(),
882            Life {
883                tokens,
884                created: now,
885                expires_at: lease.expires_at,
886                ended: None,
887                last_read: now,
888                line_gap: previous.longest_gap,
889            },
890        );
891        book.events.push(CacheEvent::Created {
892            name: name.clone(),
893            tokens,
894            estimated: create.expected,
895            covers: create.covers,
896            ttl_secs: create.ttl_secs,
897            at: now,
898        });
899        book.report.created.push(CreatedCache {
900            name,
901            tokens,
902            at: now,
903        });
904        book.leases.insert(lease.digest.clone(), lease.clone());
905        tracing::info!(target: "gemini.cache.created", name = %lease.name, tokens, covers = create.covers);
906        lease
907    }
908
909    pub(crate) fn create_failed(
910        &self,
911        line: &str,
912        coverable: u64,
913        status: Option<u16>,
914        message: String,
915    ) {
916        let now = self.now();
917        let mut book = self.book();
918        if let Some(entry) = book.lines.get_mut(line) {
919            entry.failed_at = Some(coverable);
920        }
921        book.events.push(CacheEvent::CreateFailed {
922            status,
923            message,
924            at: now,
925        });
926    }
927
928    pub(crate) fn extended(&self, lease: &Lease, ttl_secs: u64) {
929        let now = self.now();
930        let expires_at = now + ttl_secs;
931        let mut book = self.book();
932        if let Some(held) = book.leases.get_mut(&lease.digest) {
933            held.expires_at = expires_at;
934        }
935        if let Some(life) = book.lives.get_mut(&lease.name) {
936            life.expires_at = expires_at;
937        }
938        book.events.push(CacheEvent::Extended {
939            name: lease.name.clone(),
940            expires_at,
941            at: now,
942        });
943        tracing::info!(target: "gemini.cache.extended", name = %lease.name, expires_at);
944    }
945
946    pub(crate) fn touched(&self, lease: &Lease) {
947        let now = self.now();
948        let mut book = self.book();
949        let gap = book
950            .lines
951            .get(&lease.digest)
952            .map_or(0, |line| line.longest_gap);
953        if let Some(life) = book.lives.get_mut(&lease.name) {
954            life.last_read = now;
955            life.line_gap = life.line_gap.max(gap);
956        }
957    }
958}