agentplane/model/gemini.rs
1//! A `ModelProvider` for the Google Gemini Developer API.
2//!
3//! Behind the `gemini` feature. Speaks `generateContent` and
4//! `streamGenerateContent` directly rather than going through Google's
5//! OpenAI-compatibility endpoint, and the reason is the one thing this wire has
6//! that the compatible one cannot carry cleanly.
7//!
8//! # Why not the compatibility endpoint
9//!
10//! Gemini is reachable as `.../v1beta/openai/chat/completions`, so the
11//! `chat-completions` driver *does* reach it. What it cannot reach is a
12//! first-class contract:
13//!
14//! * **Reasoning is refused there.** The compatible wire has no
15//! model-family-neutral spelling for it, so a declared `reasoning_effort`
16//! is refused rather than silently dropped — which means the control a
17//! manifest declares does not apply to Gemini at all.
18//! * **Governed media is refused there**, for the same reason: multimodal
19//! content is a per-server dialect on that wire.
20//! * **Structured output falls back to a forced tool**, because whether a
21//! compatible server honours `json_schema` is exactly what cannot be
22//! assumed. Gemini enforces a schema natively during generation.
23//!
24//! And one thing that is not a matter of degree. Gemini's thinking models
25//! attach an encrypted **thought signature** to the parts they emit, and
26//! **reject** a follow-up turn that does not carry it back. Through the
27//! compatible endpoint it rides in
28//! `tool_calls[].extra_content.google.thought_signature` — a place a driver
29//! normalising every provider into one shape has nowhere to keep. The
30//! ecosystem has the scars: `LiteLLM` ended up smuggling the signature inside the
31//! tool-call *id*, which then leaked into requests to other providers and still
32//! degenerates multi-turn tool calling when the signature arrives on a thought
33//! part rather than a function-call part.
34//!
35//! This driver avoids that class of bug by never being in a position to have
36//! it: the continuation is the model's `content` **verbatim**, opaque to the
37//! runtime, exactly as the `OpenAI` driver carries encrypted reasoning items and
38//! the Anthropic driver carries signed thinking blocks. There is no field to
39//! know about, so there is no field to lose.
40//!
41//! # Which surface, and why this one
42//!
43//! Google's **Interactions API** is generally available and is what it
44//! recommends for new work. This driver deliberately targets `generateContent`,
45//! and the reason is an invariant rather than inertia: the Interactions API's
46//! defining feature is server-side conversation state addressed by
47//! `previous_interaction_id`, retained for 55 days on the paid tier and **one
48//! day** on the free one. Provider-held conversation state cannot be replay
49//! truth here — a run replayed after the retention window would have nothing to
50//! replay against. A deployment would therefore have to set `store=false`,
51//! which removes the API's own advantage and, at the time of writing, is the
52//! mode whose thought-signature handling Google does not document.
53//!
54//! `generateContent` is stateless by construction, carries the signature
55//! explicitly for the client to return, and remains fully supported. Revisit
56//! when the Interactions API documents stateless multi-turn function calling —
57//! not merely when it gains features, because features behind server-side state
58//! are features this runtime cannot use.
59//!
60//! # What this driver deliberately does not send
61//!
62//! `GenerationConfig` has knobs this driver never touches, and the omissions
63//! are decisions rather than gaps:
64//!
65//! * **`temperature`, `topP`, `topK`, `seed`** — no sampling parameter is sent
66//! by any driver here. They are absent from [`Request`], so they could not
67//! enter the effect key, and a knob that changes what the provider does
68//! without changing effect identity is one a replay cannot account for.
69//! `seed` is the tempting one and is the clearest case: replay here never
70//! calls the model again, so a seed buys nothing, and sending one would imply
71//! a reproducibility guarantee no provider actually makes.
72//! * **`candidateCount`** — left at its default of one. The runtime journals
73//! one completion per effect; asking for several answers and picking one
74//! would make the choice an unrecorded step.
75//! * **`responseModalities`, `speechConfig`, `mediaResolution`** — image and
76//! audio *output*. A completion is text plus tool calls; generated bytes
77//! would need blob storage, a digest in the chain and a retention unit, which
78//! is a feature rather than a field.
79//! * **`stopSequences`** — no seam declares them.
80//!
81//! Two fields are chosen rather than merely used. Schemas go in
82//! **`responseJsonSchema`**, not `responseSchema`: the latter is Gemini's
83//! trimmed OpenAPI-subset dialect, and translating a caller's JSON Schema into
84//! it would mean the effect key records one shape while the wire carries
85//! another — the quiet divergence this crate refuses everywhere. And thinking
86//! is requested as **`thinkingLevel`**, the form the Gemini 3 models take,
87//! rather than the token-denominated `thinkingBudget` of the 2.5 generation: a
88//! level is what [`ReasoningEffort`] means, and converting one into a token
89//! count would be this driver inventing a number nobody declared.
90//!
91//! # The failure table
92//!
93//! Status classification is shared doctrine, common to every HTTP driver here.
94//! What is
95//! specific here is the success envelope: a `finishReason` of `MAX_TOKENS` is a
96//! truncated answer reported through [`Completion::truncated`] rather than as a
97//! silently shortened string; `STOP` is the only other reason read as an answer,
98//! and every other one — `SAFETY`, `OTHER`, `MALFORMED_FUNCTION_CALL`, a reason
99//! not yet invented — is a **metered** decline, because deciding to stop cost
100//! whatever it cost;
101//! and a candidate with no parts at all is a loud `Unusable` rather than an
102//! empty answer.
103
104use async_trait::async_trait;
105use serde_json::{Value, json};
106
107use crate::core::Secret;
108
109use super::wire::{RESPOND_TOOL, classify_status, classify_transport, structured};
110use super::{
111 Completion, ModelError, ModelId, ModelProvider, ReasoningEffort, Request, SchemaMode, Usage,
112 gemini_stream, sse,
113};
114
115/// The provider tag a continuation from this driver carries.
116pub(crate) const PROVIDER: &str = "gemini";
117
118/// A harm category Gemini can be asked to block.
119///
120/// Named as an enum rather than taken as a string for the reason every other
121/// declared control here is typed: a misspelled category is not an error on
122/// this wire, it is a setting that silently governs nothing, and the deployment
123/// that wrote it would believe the opposite.
124#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
125#[non_exhaustive]
126pub enum HarmCategory {
127 Harassment,
128 HateSpeech,
129 SexuallyExplicit,
130 DangerousContent,
131 CivicIntegrity,
132 Jailbreak,
133}
134
135impl HarmCategory {
136 /// The wire spelling.
137 #[must_use]
138 pub const fn as_str(self) -> &'static str {
139 match self {
140 Self::Harassment => "HARM_CATEGORY_HARASSMENT",
141 Self::HateSpeech => "HARM_CATEGORY_HATE_SPEECH",
142 Self::SexuallyExplicit => "HARM_CATEGORY_SEXUALLY_EXPLICIT",
143 Self::DangerousContent => "HARM_CATEGORY_DANGEROUS_CONTENT",
144 Self::CivicIntegrity => "HARM_CATEGORY_CIVIC_INTEGRITY",
145 Self::Jailbreak => "HARM_CATEGORY_JAILBREAK",
146 }
147 }
148}
149
150/// How much of a category to block.
151#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
152#[non_exhaustive]
153pub enum HarmBlockThreshold {
154 LowAndAbove,
155 MediumAndAbove,
156 OnlyHigh,
157 /// Block nothing in this category.
158 ///
159 /// Present because a deployment may genuinely need it — a content-moderation
160 /// agent has to be able to read what it moderates — and because expressing
161 /// it *explicitly* is better than the alternative. An absent setting is the
162 /// provider's default, which Google may change; `None` is a decision, and it
163 /// is in the request profile where a reviewer and a replay can both see it.
164 None,
165}
166
167impl HarmBlockThreshold {
168 /// The wire spelling.
169 #[must_use]
170 pub const fn as_str(self) -> &'static str {
171 match self {
172 Self::LowAndAbove => "BLOCK_LOW_AND_ABOVE",
173 Self::MediumAndAbove => "BLOCK_MEDIUM_AND_ABOVE",
174 Self::OnlyHigh => "BLOCK_ONLY_HIGH",
175 Self::None => "BLOCK_NONE",
176 }
177 }
178}
179
180/// The deployment's own safety thresholds, passed through.
181///
182/// The same posture as the Bedrock guardrail and for the same reasons. This
183/// crate ships no content classifier — a deployment that needs one already has
184/// a better one, administered where its compliance people can see it — and what
185/// the runtime owns is everything *around* it:
186///
187/// * the settings are **effect identity**, so loosening a threshold between a
188/// run and its replay is divergence rather than a silent change in what
189/// governed the call;
190/// * an intervention is a **metered refusal**, never an answer: a prompt
191/// blocked before generation is a `Refused` naming its reason, and a
192/// generation stopped by `SAFETY` is `Unusable` carrying the tokens it
193/// burned, because deciding to stop cost whatever it cost;
194/// * both request paths carry it, since a control the streaming path drops is
195/// one a `stream: true` deployment loses.
196///
197/// Deliberately not a default. An empty set means *the provider's defaults*,
198/// which is what a deployment gets that configures none, and inventing a
199/// house policy here would be this crate deciding a question it has no standing
200/// to decide.
201#[derive(Debug, Clone, Default, PartialEq, Eq)]
202pub struct SafetySettings {
203 /// Ordered by category, so two deployments that configured the same
204 /// thresholds in a different order produce the same request bytes — and
205 /// therefore the same effect identity, rather than a spurious divergence.
206 thresholds: std::collections::BTreeMap<HarmCategory, HarmBlockThreshold>,
207}
208
209impl SafetySettings {
210 #[must_use]
211 pub fn new() -> Self {
212 Self::default()
213 }
214
215 /// Set one category's threshold, replacing any previous one.
216 #[must_use]
217 pub fn block(mut self, category: HarmCategory, threshold: HarmBlockThreshold) -> Self {
218 self.thresholds.insert(category, threshold);
219 self
220 }
221
222 #[must_use]
223 pub fn is_empty(&self) -> bool {
224 self.thresholds.is_empty()
225 }
226
227 /// The `safetySettings` array.
228 fn wire(&self) -> Value {
229 Value::Array(
230 self.thresholds
231 .iter()
232 .map(|(category, threshold)| {
233 json!({ "category": category.as_str(), "threshold": threshold.as_str() })
234 })
235 .collect(),
236 )
237 }
238
239 /// The form that enters the request profile.
240 ///
241 /// The same pairs the wire carries. A profile that recorded only *whether*
242 /// safety was configured would let a threshold move from
243 /// `BLOCK_LOW_AND_ABOVE` to `BLOCK_NONE` without changing effect identity,
244 /// which is precisely the change worth catching.
245 fn profile(&self) -> Value {
246 self.wire()
247 }
248}
249
250/// Calls the Gemini Developer API.
251///
252/// The key is held here and never journaled — transport metadata in exactly the
253/// sense a peer credential is, and the same rule applies: a secret in a
254/// hash-chained record cannot be redacted afterwards, only discovered.
255pub struct Gemini {
256 http: reqwest::Client,
257 key: Secret,
258 base: String,
259 version: String,
260 default_schema_mode: SchemaMode,
261 schema_modes: std::collections::BTreeMap<String, SchemaMode>,
262 stream: bool,
263 egress: Option<crate::core::Egress>,
264 timeout: std::time::Duration,
265 /// The deployment's own safety thresholds, if it declared any.
266 safety: SafetySettings,
267}
268
269impl std::fmt::Debug for Gemini {
270 /// Redacts the key. Deriving `Debug` would print it into every log line and
271 /// span that touches the provider.
272 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
273 f.debug_struct("Gemini")
274 .field("base", &self.base)
275 .field("version", &self.version)
276 .field("key", &"<redacted>")
277 .finish_non_exhaustive()
278 }
279}
280
281impl Gemini {
282 pub const DEFAULT_TIMEOUT: std::time::Duration = std::time::Duration::from_mins(5);
283
284 /// The default endpoint.
285 pub const DEFAULT_BASE: &'static str = "https://generativelanguage.googleapis.com";
286
287 /// The API version this driver speaks.
288 ///
289 /// Pinned rather than tracking the newest, for the reason every version
290 /// here is pinned: a provider that changes its response shape under a
291 /// running plane changes what replay reads back.
292 pub const VERSION: &'static str = "v1beta";
293
294 /// Build with a key.
295 ///
296 /// # Errors
297 ///
298 /// If the HTTP client cannot be built.
299 pub fn new(key: impl Into<String>) -> Result<Self, ModelError> {
300 let http = crate::netguard::guarded_client(crate::netguard::Reach::Configured)
301 .build()
302 .map_err(|e| ModelError::Unreachable {
303 model: ModelId::new(PROVIDER, "*"),
304 detail: format!("could not build an HTTP client: {e}"),
305 })?;
306 Ok(Self {
307 http,
308 key: Secret::new(key),
309 base: Self::DEFAULT_BASE.to_owned(),
310 version: Self::VERSION.to_owned(),
311 // Native: Gemini enforces a response schema during generation, so
312 // the forced tool is a fallback rather than the honest default it
313 // has to be on a wire whose servers only imitate the shape.
314 default_schema_mode: SchemaMode::Native,
315 schema_modes: std::collections::BTreeMap::new(),
316 stream: true,
317 egress: None,
318 timeout: Self::DEFAULT_TIMEOUT,
319 safety: SafetySettings::new(),
320 })
321 }
322
323 /// Take the key from `GEMINI_API_KEY`, falling back to `GOOGLE_API_KEY`.
324 ///
325 /// Both are read because both are in wide use and a deployment that
326 /// exported the other one would otherwise get an authentication failure
327 /// naming neither.
328 ///
329 /// # Errors
330 ///
331 /// If neither variable is set, or the HTTP client cannot be built.
332 pub fn from_env() -> Result<Self, ModelError> {
333 let key = std::env::var("GEMINI_API_KEY")
334 .or_else(|_| std::env::var("GOOGLE_API_KEY"))
335 .map_err(|_| ModelError::Refused {
336 model: ModelId::new(PROVIDER, "*"),
337 detail: "neither GEMINI_API_KEY nor GOOGLE_API_KEY is set".to_owned(),
338 })?;
339 Self::new(key)
340 }
341
342 /// Point at another endpoint — a regional host, or a test double.
343 #[must_use]
344 pub fn base(mut self, base: impl Into<String>) -> Self {
345 let mut base = base.into();
346 while base.ends_with('/') {
347 base.pop();
348 }
349 self.base = base;
350 self
351 }
352
353 /// Bound connection, generation, and response streaming as one operation.
354 #[must_use]
355 pub const fn timeout(mut self, timeout: std::time::Duration) -> Self {
356 self.timeout = timeout;
357 self
358 }
359
360 /// How to obtain a schema-conforming answer from every model.
361 #[must_use]
362 pub fn structured_via(mut self, mode: SchemaMode) -> Self {
363 self.default_schema_mode = mode;
364 self
365 }
366
367 /// How to obtain a schema-conforming answer from **one** model.
368 #[must_use]
369 pub fn structured_via_for(mut self, model: impl Into<String>, mode: SchemaMode) -> Self {
370 self.schema_modes.insert(model.into(), mode);
371 self
372 }
373
374 /// Apply the deployment's own safety thresholds to every call.
375 ///
376 /// See [`SafetySettings`] for the posture. In short: this crate ships no
377 /// classifier, Google's is configured here, and what the runtime owns is
378 /// that the choice is **effect identity** and that an intervention is a
379 /// metered refusal rather than an answer.
380 #[must_use]
381 pub fn safety(mut self, safety: SafetySettings) -> Self {
382 self.safety = safety;
383 self
384 }
385
386 /// Restrict where this driver may connect.
387 ///
388 /// Deny-by-default once set — see [`Egress`](crate::core::Egress).
389 #[must_use]
390 pub fn egress(mut self, egress: crate::core::Egress) -> Self {
391 self.egress = Some(egress);
392 self
393 }
394
395 /// Ask for the whole response at once instead of streaming it.
396 ///
397 /// Streaming is the default for the reason it is elsewhere in this crate:
398 /// Gemini reports `usageMetadata` in the stream, so a severed connection can
399 /// say what it burned rather than being billed as zero and retried for free.
400 #[must_use]
401 pub const fn buffered(mut self) -> Self {
402 self.stream = false;
403 self
404 }
405
406 fn mode_for(&self, model: &ModelId) -> SchemaMode {
407 self.schema_modes
408 .get(&model.model)
409 .copied()
410 .unwrap_or(self.default_schema_mode)
411 }
412
413 fn check_egress(&self, model: &ModelId) -> Result<(), ModelError> {
414 let Some(egress) = &self.egress else {
415 return Ok(());
416 };
417 let host = reqwest::Url::parse(&self.base)
418 .ok()
419 .and_then(|u| u.host_str().map(ToOwned::to_owned));
420 egress
421 .permits(host.as_deref())
422 .map_err(|e| ModelError::Egress {
423 model: model.clone(),
424 detail: e.to_string(),
425 })
426 }
427
428 fn refused(model: &ModelId, detail: impl Into<String>) -> ModelError {
429 ModelError::Refused {
430 model: model.clone(),
431 detail: detail.into(),
432 }
433 }
434
435 /// `thinkingConfig`, or why this effort cannot be asked for.
436 ///
437 /// Gemini names four levels — `minimal`, `low`, `medium`, `high` — and this
438 /// maps the four that match exactly. The rest are **refused rather than
439 /// collapsed**: `reasoning_effort` is digest-covered precisely so it
440 /// describes what governed a call, and answering a request for `max` with
441 /// the highest level that happens to exist is a substitution nothing
442 /// downstream could see. The same rule the Anthropic and Bedrock drivers
443 /// follow for the levels their providers cannot express.
444 ///
445 /// `None` is refused rather than mapped to `minimal`, and the distinction
446 /// matters: Google documents that thinking **cannot be turned off** on the
447 /// Gemini 3 models, so a driver rendering "do not reason" as "reason a
448 /// little" would report a control as applied that the provider never
449 /// applied.
450 fn thinking_config(model: &ModelId, effort: ReasoningEffort) -> Result<Value, ModelError> {
451 let level = match effort {
452 ReasoningEffort::Minimal => "minimal",
453 ReasoningEffort::Low => "low",
454 ReasoningEffort::Medium => "medium",
455 ReasoningEffort::High => "high",
456 ReasoningEffort::None | ReasoningEffort::XHigh | ReasoningEffort::Max => {
457 return Err(Self::refused(
458 model,
459 format!(
460 "Gemini has no thinking level for reasoning effort '{}' — it names \
461 minimal, low, medium and high, and thinking cannot be switched off \
462 on the Gemini 3 models",
463 effort.as_str()
464 ),
465 ));
466 }
467 };
468 Ok(json!({ "thinkingLevel": level }))
469 }
470
471 /// The `contents` array, from whatever shape the caller supplied.
472 ///
473 /// A bare string is one user turn. An array is already `contents` and is
474 /// passed through untouched — which is how governed media reaches this
475 /// driver, as `inlineData` parts the caller assembled and the runtime
476 /// materialised. An object is read for `messages` only when
477 /// `prompt_envelope` says it is an envelope rather than content; anything
478 /// else is one user turn carrying the object.
479 fn contents(prompt: &Value) -> Value {
480 match prompt {
481 Value::String(text) => json!([{ "role": "user", "parts": [{ "text": text }] }]),
482 Value::Array(_) => prompt.clone(),
483 other => crate::model::prompt_envelope(other, "messages")
484 .cloned()
485 .unwrap_or_else(|| {
486 // `system` is an instruction *about* the content and is lifted
487 // out separately; leaving it here would show the model its own
488 // orders as part of the question.
489 let mut rest = other.clone();
490 if let Some(map) = rest.as_object_mut() {
491 map.remove("system");
492 }
493 json!([{ "role": "user", "parts": [{ "text": rest.to_string() }] }])
494 }),
495 }
496 }
497
498 /// The system instruction, if the caller set one.
499 ///
500 /// Gemini takes this as a **top-level** `systemInstruction`, not a role
501 /// inside `contents`. A driver that left it in the turns would show the
502 /// model an instruction it treats as ordinary content — and the instruction
503 /// slot is authority-bearing here, so that is not a cosmetic difference.
504 fn system_instruction(prompt: &Value) -> Option<Value> {
505 let system = prompt.get("system").filter(|s| !s.is_null())?;
506 Some(match system {
507 Value::String(text) => json!({ "parts": [{ "text": text }] }),
508 // Already a `Content` — passed through, so a caller can supply
509 // several parts.
510 other => other.clone(),
511 })
512 }
513
514 /// The request body.
515 fn body(&self, model: &ModelId, request: &Request<'_>) -> Result<Value, ModelError> {
516 let Request {
517 prompt,
518 max_output_tokens,
519 reasoning_effort,
520 schema,
521 tools,
522 exchanges,
523 continuation,
524 ..
525 } = request;
526
527 let mut contents = Self::contents(prompt);
528 Self::append_tool_turns(&mut contents, exchanges, *continuation, model)?;
529
530 let mut generation_config = json!({ "maxOutputTokens": max_output_tokens });
531 if let Some(effort) = reasoning_effort {
532 generation_config["thinkingConfig"] = Self::thinking_config(model, *effort)?;
533 }
534
535 let mut body = json!({ "contents": contents });
536 // Both request paths get it from here, because both build their body
537 // here — the Bedrock guardrail had to say this twice and that is the
538 // shape where only the half nobody exercises is wrong.
539 if !self.safety.is_empty() {
540 body["safetySettings"] = self.safety.wire();
541 }
542 if let Some(system) = Self::system_instruction(prompt) {
543 body["systemInstruction"] = system;
544 }
545
546 let mode = self.mode_for(model);
547 let mut declarations: Vec<Value> = tools
548 .iter()
549 .map(|t| {
550 json!({
551 "name": t.name,
552 "description": t.description,
553 // `parametersJsonSchema` rather than `parameters`: the
554 // latter is Gemini's own trimmed Schema dialect, and
555 // handing it a full JSON Schema is how a valid declaration
556 // becomes a 400 nobody can read.
557 "parametersJsonSchema": t.parameters,
558 })
559 })
560 .collect();
561
562 if let Some(schema) = schema {
563 match mode {
564 SchemaMode::Native => {
565 generation_config["responseMimeType"] = json!("application/json");
566 generation_config["responseJsonSchema"] = (*schema).clone();
567 }
568 SchemaMode::ForcedTool => {
569 if !tools.is_empty() {
570 return Err(Self::refused(
571 model,
572 "forced-tool structured output cannot be combined with declared \
573 tools: the model would be offered a choice between answering and \
574 calling one. Use SchemaMode::Native, which Gemini enforces during \
575 generation",
576 ));
577 }
578 declarations.push(json!({
579 "name": RESPOND_TOOL,
580 "description": "Return the answer in the required shape.",
581 "parametersJsonSchema": (*schema).clone(),
582 }));
583 body["toolConfig"] = json!({
584 "functionCallingConfig": {
585 "mode": "ANY",
586 "allowedFunctionNames": [RESPOND_TOOL],
587 }
588 });
589 }
590 }
591 }
592
593 if !declarations.is_empty() {
594 body["tools"] = json!([{ "functionDeclarations": declarations }]);
595 }
596 body["generationConfig"] = generation_config;
597 Ok(body)
598 }
599
600 /// Append the model turn that asked for tools, and the results.
601 ///
602 /// The model turn is the continuation **verbatim**. That is the whole
603 /// reason this driver exists rather than deferring to the compatible wire:
604 /// a part may carry a `thoughtSignature`, Gemini rejects a follow-up turn
605 /// that does not return it, and a driver reconstructing the turn from the
606 /// fields it understands cannot return what it never kept.
607 fn append_tool_turns(
608 contents: &mut Value,
609 exchanges: &[super::ToolExchange],
610 continuation: Option<&super::ProviderContinuation>,
611 model: &ModelId,
612 ) -> Result<(), ModelError> {
613 if exchanges.is_empty() {
614 if continuation.is_some() {
615 // Silently dropping it would journal an effect key that
616 // records a continuation the wire never carried.
617 return Err(Self::refused(
618 model,
619 "a continuation without tool exchanges has no request to follow",
620 ));
621 }
622 return Ok(());
623 }
624 let Some(array) = contents.as_array_mut() else {
625 return Err(Self::refused(
626 model,
627 "the prompt did not assemble into a `contents` array",
628 ));
629 };
630
631 match continuation {
632 Some(state) if state.provider == PROVIDER => match state.state.as_array() {
633 Some(turns) => array.extend(turns.iter().cloned()),
634 None => {
635 return Err(Self::refused(
636 model,
637 "the continuation was not a Gemini contents array",
638 ));
639 }
640 },
641 Some(other) => {
642 return Err(Self::refused(
643 model,
644 format!(
645 "the continuation was issued by '{}' and this is the Gemini driver — \
646 provider state is opaque and is never valid across providers",
647 other.provider
648 ),
649 ));
650 }
651 // No continuation: a caller assembled the exchanges by hand. The
652 // turn is rebuilt, and it is worth naming what that costs, because
653 // it is the failure this driver exists to avoid — a rebuilt turn
654 // carries no thought signature, and a thinking model will refuse
655 // it. Callers who ran the model through this driver always have
656 // one.
657 None => array.push(json!({
658 "role": "model",
659 "parts": exchanges
660 .iter()
661 .map(|e| json!({
662 "functionCall": { "name": e.call.name, "args": e.call.arguments }
663 }))
664 .collect::<Vec<_>>(),
665 })),
666 }
667
668 array.push(Self::tool_responses(exchanges));
669 Ok(())
670 }
671
672 /// The results turn: one user content of `functionResponse` parts, in the
673 /// order the calls were made.
674 fn tool_responses(exchanges: &[super::ToolExchange]) -> Value {
675 json!({
676 "role": "user",
677 "parts": exchanges
678 .iter()
679 .map(|e| {
680 let mut response = json!({ "name": e.call.name, "response": {
681 // Gemini expects an object; a bare value is wrapped
682 // rather than sent as-is, because a string here is a
683 // 400 rather than a smaller answer.
684 "output": e.output,
685 }});
686 if e.failed {
687 response["response"] = json!({ "error": e.output });
688 }
689 // A call the provider identified is answered under its id;
690 // an id this driver synthesized was never the provider's.
691 if !is_synthesized_id(&e.call.name, &e.call.id) {
692 response["id"] = json!(e.call.id);
693 }
694 json!({ "functionResponse": response })
695 })
696 .collect::<Vec<_>>(),
697 })
698 }
699
700 /// Extend the provider transcript after a successful tool-calling
701 /// response, exactly as every sibling driver does: prior rounds, then this
702 /// round's results, then the turn just answered. The runtime clears
703 /// `exchanges` each turn and relies on this being the whole history.
704 fn accumulate_continuation(
705 completion: &mut Completion,
706 prior: Option<&super::ProviderContinuation>,
707 exchanges: &[super::ToolExchange],
708 ) {
709 let Some(current) = completion.continuation.as_mut() else {
710 return;
711 };
712 let mut state = prior
713 .and_then(|value| value.state.as_array())
714 .cloned()
715 .unwrap_or_default();
716 if !exchanges.is_empty() {
717 state.push(Self::tool_responses(exchanges));
718 }
719 if let Some(turns) = current.state.as_array() {
720 state.extend(turns.iter().cloned());
721 }
722 current.state = Value::Array(state);
723 }
724
725 /// Turn a response envelope into a [`Completion`], or say why it is not one.
726 ///
727 /// Shared by the buffered and streaming paths, so the two cannot disagree
728 /// about what a usable answer is.
729 fn interpret(
730 &self,
731 parsed: &Value,
732 model: &ModelId,
733 schema: Option<&Value>,
734 ) -> Result<Completion, ModelError> {
735 let usage = Self::usage(parsed);
736
737 let Some(candidate) = parsed.get("candidates").and_then(|c| c.get(0)) else {
738 // No candidate at all: either the prompt was blocked before
739 // generating, which the feedback names, or the response is one this
740 // driver cannot read. Both are refusals; only the first can say why.
741 let detail = parsed
742 .get("promptFeedback")
743 .and_then(|f| f.get("blockReason"))
744 .and_then(Value::as_str)
745 .map_or_else(
746 || "the response carried no candidates".to_owned(),
747 |reason| format!("the prompt was blocked before generating: {reason}"),
748 );
749 return Err(Self::refused(model, detail));
750 };
751
752 let finish = candidate
753 .get("finishReason")
754 .and_then(Value::as_str)
755 .map(ToOwned::to_owned);
756 // An allowlist: `STOP` is an answer, `MAX_TOKENS` a typed truncation,
757 // and every other reason — a filter, `OTHER`, a malformed or
758 // unexpected tool call, a reason this driver has never heard of — is
759 // generation that ended without one. Metered, because deciding cost
760 // whatever it cost. A denylist here passes the next reason Google adds
761 // as a complete answer.
762 let truncated = match finish.as_deref() {
763 Some("STOP") => false,
764 Some("MAX_TOKENS") => true,
765 other => {
766 return Err(ModelError::Unusable {
767 model: model.clone(),
768 usage,
769 detail: format!(
770 "generation stopped: {}",
771 other.unwrap_or("no finishReason, so completeness is unknown")
772 ),
773 });
774 }
775 };
776
777 let content = candidate.get("content").cloned().unwrap_or(Value::Null);
778 let parts = content
779 .get("parts")
780 .and_then(Value::as_array)
781 .cloned()
782 .unwrap_or_default();
783
784 let Scanned {
785 text,
786 calls,
787 forced,
788 } = Self::scan(&parts);
789
790 let emulating = schema.is_some() && self.mode_for(model) == SchemaMode::ForcedTool;
791 if text.is_empty() && calls.is_empty() && forced.is_none() && !truncated {
792 return Err(ModelError::Unusable {
793 model: model.clone(),
794 usage,
795 detail: format!("the answer carried no content (finishReason {finish:?})"),
796 });
797 }
798
799 let (text, structured_value) = if emulating {
800 let Some(arguments) = forced else {
801 return Err(ModelError::Unusable {
802 model: model.clone(),
803 usage,
804 detail: "a tool call was forced and the answer carried none — the model \
805 did not honour the function-calling config"
806 .to_owned(),
807 });
808 };
809 let raw = arguments.to_string();
810 let parsed_schema = structured(schema, &raw, &[], model, usage)?;
811 (raw, parsed_schema)
812 } else {
813 let parsed_schema = structured(schema, &text, &calls, model, usage)?;
814 (text, parsed_schema)
815 };
816
817 // The model's own turn, byte for byte, so the next request returns
818 // every signature exactly where it sat. An array of contents, not the
819 // bare turn: the state must carry *every* prior round, or round
820 // three's request forgets round one's signed turn — silently, since
821 // the model simply re-asks with amnesia.
822 let continuation = (!calls.is_empty() && content.is_object())
823 .then(|| super::ProviderContinuation::new(PROVIDER, json!([content.clone()])));
824
825 Ok(Completion {
826 structured: structured_value,
827 tool_calls: calls,
828 text,
829 // Reported at the top of the response rather than on the candidate,
830 // and on every chunk of a stream — so the accumulator's envelope
831 // carries it and one read covers both paths.
832 model: parsed
833 .get("modelVersion")
834 .and_then(Value::as_str)
835 .map(ToOwned::to_owned),
836 usage,
837 stop_reason: finish,
838 truncated,
839 continuation,
840 })
841 }
842
843 /// Usage, in this crate's terms.
844 ///
845 /// Two normalisations, both of which cost real money when got wrong.
846 /// `thoughtsTokenCount` is billed as **output** and is reported *beside*
847 /// `candidatesTokenCount` rather than inside it, so it is added — a
848 /// reasoning-heavy run would otherwise under-report its bill by most of it.
849 /// `cachedContentTokenCount` is a **subset** of `promptTokenCount`, as
850 /// `OpenAI` reports cached input, so it is recorded rather than added.
851 fn usage(parsed: &Value) -> Usage {
852 let count = |key: &str| {
853 parsed
854 .get("usageMetadata")
855 .and_then(|u| u.get(key))
856 .and_then(Value::as_u64)
857 .unwrap_or_default()
858 };
859 Usage {
860 input_tokens: count("promptTokenCount"),
861 // Saturating: both counts are whatever the response said, and a
862 // wrapped sum reads an astronomical bill as a free one.
863 output_tokens: count("candidatesTokenCount")
864 .saturating_add(count("thoughtsTokenCount")),
865 cache_read_tokens: count("cachedContentTokenCount"),
866 cache_write_tokens: 0,
867 minor_units: 0,
868 }
869 }
870
871 fn url(&self, model: &ModelId) -> String {
872 let method = if self.stream {
873 "streamGenerateContent?alt=sse"
874 } else {
875 "generateContent"
876 };
877 format!(
878 "{}/{}/models/{}:{method}",
879 self.base, self.version, model.model
880 )
881 }
882
883 async fn read_buffered(
884 &self,
885 response: reqwest::Response,
886 model: &ModelId,
887 schema: Option<&Value>,
888 ) -> Result<Completion, ModelError> {
889 // Read under this plane's ceiling rather than to end-of-stream: a
890 // provider is a counterparty, and a counterparty must not decide how
891 // much of this process's memory its answer costs.
892 let body = crate::netguard::intake::read(response, crate::netguard::intake::ANSWER)
893 .await
894 .map_err(|e| super::wire::classify_intake(model, Usage::default(), &e))?;
895 let parsed: Value = serde_json::from_slice(&body).map_err(|e| ModelError::Unusable {
896 model: model.clone(),
897 usage: Usage::default(),
898 detail: format!("the response did not parse: {e}"),
899 })?;
900 self.interpret(&parsed, model, schema)
901 }
902
903 async fn read_streamed(
904 &self,
905 response: reqwest::Response,
906 model: &ModelId,
907 schema: Option<&Value>,
908 observer: Option<(&dyn super::ModelStreamObserver, &crate::core::Label)>,
909 ) -> Result<Completion, ModelError> {
910 use futures_util::StreamExt;
911
912 let mut decoder = sse::Decoder::new();
913 let mut acc = gemini_stream::Accumulator::new();
914 let mut body = response.bytes_stream();
915 // The same ceiling the buffered path applies, to the same bytes.
916 // `sse::Decoder` already bounds one event, which is the unterminated
917 // line; this bounds the *number* of them. A stream of well-formed
918 // hundred-byte deltas passes every check the decoder makes and grows
919 // the accumulator until the process dies.
920 let mut meter = crate::netguard::intake::Meter::new(crate::netguard::intake::ANSWER);
921
922 while let Some(chunk) = body.next().await {
923 let chunk = match chunk {
924 Ok(chunk) => chunk,
925 Err(e) => {
926 return Err(severed(model, &acc, &crate::netguard::transport_text(&e)));
927 }
928 };
929 // Charged before the chunk is kept: what the ceiling bounds is
930 // what this process holds, not what it has already held. The
931 // refusal is `Unusable` rather than `severed` — it is this
932 // plane's rather than the provider's, and the call generated.
933 // It carries the cumulative usage the chunks already reported.
934 if let Err(e) = meter.charge(chunk.len()) {
935 let usage = acc
936 .usage_envelope()
937 .map_or_else(Usage::default, |envelope| Self::usage(&envelope));
938 return Err(super::wire::classify_intake(model, usage, &e));
939 }
940 let events = decoder
941 .push(&chunk)
942 .map_err(|error| severed(model, &acc, &error.to_string()))?;
943 for event in events {
944 if let Some(delta) = acc.push(&event.data)
945 && let Some((observer, label)) = observer
946 {
947 observer.event(crate::core::Tainted::with_label(
948 super::ModelStreamEvent::TextDelta(delta),
949 label.clone(),
950 ));
951 }
952 }
953 }
954 if !acc.done() {
955 // A blocked prompt is a refusal with a name, not an outage: the
956 // stream ends without a finish reason, and classifying that as
957 // severed would mark it safe to repeat — a retry loop re-hitting
958 // the same block forever.
959 if acc.prompt_blocked() {
960 return self.interpret(&acc.into_response(), model, schema);
961 }
962 return Err(severed(
963 model,
964 &acc,
965 "the stream ended before the model said why it stopped",
966 ));
967 }
968
969 let completion = self.interpret(&acc.into_response(), model, schema)?;
970 if let Some((observer, label)) = observer {
971 observer.event(crate::core::Tainted::with_label(
972 super::ModelStreamEvent::Usage(completion.usage),
973 label.clone(),
974 ));
975 }
976 Ok(completion)
977 }
978}
979
980/// What one pass over a candidate's parts found.
981struct Scanned {
982 text: String,
983 calls: Vec<super::ToolCall>,
984 /// The forced-tool answer, when structured output is being emulated.
985 forced: Option<Value>,
986}
987
988impl Gemini {
989 /// Read a candidate's parts into text, tool calls, and a forced answer.
990 ///
991 /// Split out of `interpret` so the interesting decisions are visible rather
992 /// than buried in the middle of a long function: which parts are *not* the
993 /// answer, and where a tool call's id comes from when the provider issues
994 /// none.
995 fn scan(parts: &[Value]) -> Scanned {
996 let mut text = String::new();
997 let mut calls = Vec::new();
998 let mut forced: Option<Value> = None;
999 for (index, part) in parts.iter().enumerate() {
1000 // A part marked `thought` is opaque reasoning. It stays in the
1001 // continuation and never becomes the answer: exposing it would make
1002 // the model's internal deliberation read as its conclusion.
1003 if part.get("thought").and_then(Value::as_bool) == Some(true) {
1004 continue;
1005 }
1006 if let Some(call) = part.get("functionCall") {
1007 let name = call
1008 .get("name")
1009 .and_then(Value::as_str)
1010 .unwrap_or_default()
1011 .to_owned();
1012 let arguments = call.get("args").cloned().unwrap_or_else(|| json!({}));
1013 if name == RESPOND_TOOL {
1014 forced = Some(arguments);
1015 continue;
1016 }
1017 calls.push(super::ToolCall {
1018 // Gemini's function calls need not carry an id, and the
1019 // runtime keys tool results by one. Derived from the
1020 // position rather than generated, so it is stable across a
1021 // replay of the same recorded response — a random id would
1022 // make a replayed completion differ from the one journaled.
1023 id: call
1024 .get("id")
1025 .and_then(Value::as_str)
1026 .map_or_else(|| synthesized_id(&name, index), ToOwned::to_owned),
1027 name,
1028 arguments,
1029 });
1030 continue;
1031 }
1032 if let Some(chunk) = part.get("text").and_then(Value::as_str) {
1033 text.push_str(chunk);
1034 }
1035 }
1036 Scanned {
1037 text,
1038 calls,
1039 forced,
1040 }
1041 }
1042}
1043
1044/// The id given a function call the provider sent without one: its name and
1045/// its position in the response.
1046fn synthesized_id(name: &str, index: usize) -> String {
1047 format!("{name}-{index}")
1048}
1049
1050/// Whether `id` has the shape [`synthesized_id`] gives a call named `name`.
1051fn is_synthesized_id(name: &str, id: &str) -> bool {
1052 id.strip_prefix(name)
1053 .and_then(|rest| rest.strip_prefix('-'))
1054 .is_some_and(|n| !n.is_empty() && n.bytes().all(|b| b.is_ascii_digit()))
1055}
1056
1057/// A stream that stopped before the model said why.
1058///
1059/// Three rungs, and which one applies is decided by what the wire has already
1060/// said rather than by how far the answer got:
1061///
1062/// 1. **`usageMetadata` seen** — Gemini reports it on the chunks themselves,
1063/// cumulatively, so a stream cut off mid-answer has already been told what
1064/// it burned. That is [`ModelError::Interrupted`], the one severed-stream
1065/// answer that carries a bill, and it is the reason streaming is this
1066/// driver's default.
1067/// 2. **content but no usage** — generation happened and the cost is unknown.
1068/// [`ModelError::Unaccounted`]: never free, never counted.
1069/// 3. **nothing at all** — no evidence anything reached the model, so the call
1070/// is safe to repeat and costs nothing.
1071///
1072/// Rung 1 is not an optimisation. Without it every severed Gemini stream bills
1073/// zero, and the token ceiling that exists to bound a runaway provider counts
1074/// nothing while the provider spends — the failure the ceiling was bought to
1075/// prevent, in the one driver whose wire makes it avoidable.
1076fn severed(model: &ModelId, acc: &gemini_stream::Accumulator, detail: &str) -> ModelError {
1077 // Parsed by the buffered path's own function, so the normalisation that
1078 // costs money — thought tokens billed as output, cached input a subset of
1079 // the prompt — has exactly one spelling.
1080 if let Some(envelope) = acc.usage_envelope() {
1081 return ModelError::Interrupted {
1082 model: model.clone(),
1083 usage: Gemini::usage(&envelope),
1084 detail: detail.to_owned(),
1085 };
1086 }
1087 if acc.generated() {
1088 return ModelError::Unaccounted {
1089 model: model.clone(),
1090 detail: detail.to_owned(),
1091 };
1092 }
1093 ModelError::Unavailable {
1094 model: model.clone(),
1095 detail: detail.to_owned(),
1096 }
1097}
1098
1099/// Fill a rate limit's window from the body's `google.rpc.RetryInfo`.
1100///
1101/// The shared classifier reads `Retry-After`, which is the one header every
1102/// other wire this crate meets uses — and the one this provider does not
1103/// send. Google names its window *inside* the 429 body, as a `RetryInfo`
1104/// detail carrying a proto `Duration` (`"40s"`). Discarded, the default
1105/// policy spends its attempts in milliseconds against a window measured in
1106/// tens of seconds and reports the provider down — the exact defect the
1107/// `retry_after` plumbing exists to prevent. Any other classification passes
1108/// through untouched, and a window the header already named is not replaced.
1109fn with_retry_info(error: ModelError, body: &str) -> ModelError {
1110 match error {
1111 ModelError::RateLimited {
1112 model,
1113 detail,
1114 retry_after: None,
1115 } => ModelError::RateLimited {
1116 model,
1117 detail,
1118 retry_after: retry_info_seconds(body),
1119 },
1120 other => other,
1121 }
1122}
1123
1124/// The `RetryInfo` detail's delay, in whole seconds.
1125///
1126/// Whole seconds only, floor of a fractional value, zero reads as no advice —
1127/// the same conservatisms `core::retry_after_seconds` applies to the header
1128/// form. The ceiling on believing it (`max_advice`) stays where it always
1129/// was, in the retry policy.
1130fn retry_info_seconds(body: &str) -> Option<u64> {
1131 let parsed: Value = serde_json::from_str(body).ok()?;
1132 let details = parsed.get("error")?.get("details")?.as_array()?;
1133 let delay = details
1134 .iter()
1135 .find(|d| {
1136 d.get("@type").and_then(Value::as_str)
1137 == Some("type.googleapis.com/google.rpc.RetryInfo")
1138 })?
1139 .get("retryDelay")?
1140 .as_str()?;
1141 let seconds: u64 = delay.strip_suffix('s')?.split('.').next()?.parse().ok()?;
1142 (seconds > 0).then_some(seconds)
1143}
1144
1145#[async_trait]
1146impl ModelProvider for Gemini {
1147 fn request_profile(&self, model: &ModelId) -> Value {
1148 json!({
1149 "driver": "google-gemini-generatecontent/v1",
1150 "base": self.base,
1151 "api_version": self.version,
1152 "stream": self.stream,
1153 "schema_mode": match self.mode_for(model) {
1154 SchemaMode::Native => "native",
1155 SchemaMode::ForcedTool => "forced-tool",
1156 },
1157 // Identity, not decoration: loosening a threshold changes what
1158 // governed the call, so a replay of history written under the
1159 // stricter one reports divergence rather than answering under the
1160 // looser. Absent when nothing was declared, which is a different
1161 // request from one declaring the provider defaults explicitly.
1162 "safety": (!self.safety.is_empty()).then(|| self.safety.profile()),
1163 })
1164 }
1165
1166 async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
1167 let model = request.model;
1168 super::refuse_provider_side_media(request.prompt, model)?;
1169 super::refuse_in_thread_instructions(request.prompt, model)?;
1170 self.check_egress(model)?;
1171
1172 let body = self.body(model, &request)?;
1173 let response = self
1174 .http
1175 .post(self.url(model))
1176 // The header form rather than a `?key=` query parameter: a URL is
1177 // logged by proxies, written into traces and echoed in errors, and
1178 // a credential that reaches any of those is one that cannot be
1179 // un-leaked.
1180 .header("x-goog-api-key", self.key.expose())
1181 .timeout(self.timeout)
1182 .json(&body)
1183 .send()
1184 .await
1185 .map_err(|e| classify_transport(model, &e))?;
1186
1187 let status = response.status();
1188 if !status.is_success() {
1189 let headers = response.headers().clone();
1190 // Bounded, and the ceiling is the small one: this body is read
1191 // only to say *why* the call failed, so an endpoint answering a
1192 // failure with a gigabyte gets an unexplained failure rather than
1193 // this process's memory.
1194 let text =
1195 crate::netguard::intake::read_text(response, crate::netguard::intake::METADATA)
1196 .await
1197 .unwrap_or_default();
1198 return Err(with_retry_info(
1199 classify_status(model, status.as_u16(), &headers, &text),
1200 &text,
1201 ));
1202 }
1203
1204 let mut completion = if self.stream {
1205 self.read_streamed(response, model, request.schema, request.stream)
1206 .await?
1207 } else {
1208 self.read_buffered(response, model, request.schema).await?
1209 };
1210 // A buffered call still answers the observer's one guaranteed
1211 // question — what did this cost — as every driver does on both paths.
1212 if !self.stream
1213 && let Some((observer, label)) = request.stream
1214 {
1215 observer.event(crate::core::Tainted::with_label(
1216 super::ModelStreamEvent::Usage(completion.usage),
1217 label.clone(),
1218 ));
1219 }
1220 Self::accumulate_continuation(&mut completion, request.continuation, request.exchanges);
1221 Ok(completion)
1222 }
1223}