Skip to main content

rig_core/providers/gemini/interactions_api/
mod.rs

1//! Wires for the [Gemini Interactions API](https://ai.google.dev/api/interactions-api).
2//! A request is built as JSON steps, and a reply, a whole interaction or a
3//! stream of step events, is read by [`InteractionsDecoder`](streaming::InteractionsDecoder).
4//!
5//! ```no_run
6//! use rig_core::providers::gemini::{Gemini, completion::GEMINI_2_5_FLASH};
7//!
8//! # fn main() -> Result<(), Box<dyn std::error::Error>> {
9//! let wire = Gemini::from_env()?.interactions(GEMINI_2_5_FLASH);
10//! # Ok(())
11//! # }
12//! ```
13
14use serde_json::{Map, Value, json};
15use url::form_urlencoded;
16
17use crate::completion::{CompletionRequest, Media, Replay, ReplayTarget};
18use crate::error::EncodeError;
19use crate::message::{
20    AssistantContent, DocumentMediaType, DocumentSourceKind as Source, Message, MimeType,
21    ToolChoice as Choice, ToolResultContent, UserContent,
22};
23use crate::providers::internal::wire_ids::WireIds;
24use crate::telemetry::GenAiOperation;
25use crate::wire::{Descriptor, Mode};
26
27/// Streaming helpers for the Interactions API.
28pub mod streaming;
29
30use super::completion::PROVIDER_NAME;
31
32/// The wire format both Interactions wires speak.
33const API: crate::message::Api = crate::message::Api::from_static("gemini.interactions");
34
35/// Create interactions with `POST /v1beta/interactions`.
36/// Streaming mode sets `alt=sse` and `stream: true`; unary mode reads a whole resource.
37#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
38pub struct Interactions {
39    /// The key and the API root.
40    pub provider: crate::providers::gemini::GeminiConfig,
41    /// The model to address.
42    pub model: String,
43}
44
45impl Interactions {
46    /// The wire for `model`.
47    pub fn new(provider: crate::providers::gemini::GeminiConfig, model: impl Into<String>) -> Self {
48        Self {
49            provider,
50            model: model.into(),
51        }
52    }
53}
54
55impl crate::wire::Wire for Interactions {
56    type Op = crate::operation::Completion;
57    type Payload = crate::wire::Encoded;
58    type Frame = crate::wire::WireFrame;
59    type Decoder<'id> = streaming::InteractionsDecoder;
60    type Reassembler = streaming::document::Interaction;
61
62    fn describe(&self) -> Descriptor<'_> {
63        Descriptor::new(PROVIDER_NAME)
64            .model(self.model.as_str())
65            .telemetry(|_| GenAiOperation::Chat)
66            .replay(self)
67    }
68
69    fn encode(
70        &self,
71        request: CompletionRequest,
72        mode: Mode,
73    ) -> Result<crate::wire::Encoded, EncodeError> {
74        use crate::providers::internal::LogTarget;
75        // `stream` is part of the request body on this wire, so the mode is
76        // in the bytes as well as in the path.
77        let streaming = matches!(mode, Mode::Streaming);
78        let body = create_request_body(self, &request, Some(streaming))?;
79        let (path, framing, target) = match streaming {
80            true => (
81                "/v1beta/interactions?alt=sse",
82                crate::wire::Framing::Sse,
83                LogTarget::Streaming,
84            ),
85            false => (
86                "/v1beta/interactions",
87                crate::wire::Framing::Whole,
88                LogTarget::Completions,
89            ),
90        };
91        crate::providers::internal::trace_json(
92            target,
93            "Gemini interactions completion request",
94            &body,
95        );
96        let request = http::Request::post(self.provider.interactions_uri(path))
97            .header("Content-Type", "application/json")
98            .header(
99                crate::providers::gemini::GeminiConfig::INTERACTIONS_KEY_HEADER,
100                self.provider.api_key.expose(),
101            )
102            .body(body.into_body())?;
103        // Gemini supplies no transport request-id response header.
104        Ok(crate::wire::Encoded::new(request, framing))
105    }
106
107    fn decoder<'id>(&self) -> Self::Decoder<'id> {
108        streaming::InteractionsDecoder::default()
109    }
110}
111
112impl ReplayTarget for Interactions {
113    /// Section 6.4 of the typed-options design, for Interactions.
114    fn map_options(
115        &self,
116        request: &CompletionRequest,
117        fields: crate::completion::options::OptionFields<'_>,
118    ) -> crate::completion::options::OptionMap {
119        let model = request.model.as_deref().unwrap_or(&self.model);
120        super::options::interactions(model, fields)
121    }
122
123    fn api(&self) -> crate::message::Api {
124        API
125    }
126
127    fn provider(&self) -> &str {
128        PROVIDER_NAME
129    }
130
131    fn model(&self) -> &str {
132        &self.model
133    }
134
135    /// What the model reads, by the classifier every Gemini wire shares.
136    fn accepts(&self, model: &str) -> crate::completion::Accepts {
137        super::completion::accepts(model)
138    }
139
140    fn encodes(&self, _model: &str, media: Media<'_>) -> bool {
141        encodes(media)
142    }
143
144    /// A request naming `previous_interaction_id` continues an interaction
145    /// the API stored, which holds the calls its first results answer.
146    fn continues_stored(&self, request: &CompletionRequest) -> bool {
147        crate::completion::options::param(self, request, "previous_interaction_id")
148            .is_some_and(|id| !id.is_null())
149    }
150
151    fn normalize_tool_call_id(
152        &self,
153        id: &str,
154        _model: &str,
155        _: Option<&crate::message::Origin>,
156    ) -> String {
157        crate::providers::internal::wire_ids::legal_call_id(id, 64)
158    }
159
160    /// Gemini takes system text only in `systemInstruction`: later system
161    /// messages fold into the leading one, as pi's `collapseSystemMessages`.
162    fn later_system(&self, _model: &str) -> crate::completion::LaterSystem {
163        crate::completion::LaterSystem::Leading
164    }
165
166    fn call_id_slot(&self) -> Option<&'static str> {
167        Some("/id")
168    }
169
170    /// A step's `id` survives an edit of its block.
171    fn identity(&self, item: &Value) -> Map<String, Value> {
172        item.get("id")
173            .map(|id| Map::from_iter([("id".to_owned(), id.clone())]))
174            .unwrap_or_default()
175    }
176}
177
178/// Whether this API takes `media`: data or a URL with a media type, an
179/// image of a type Gemini reads, and a document other than a PDF only as a
180/// string, which is sent as text. A file id is never taken.
181fn encodes(media: Media<'_>) -> bool {
182    let carried = |source: &Source| {
183        matches!(
184            source,
185            Source::Url(_) | Source::Base64(_) | Source::String(_)
186        )
187    };
188    match media {
189        Media::Image(image, place) => {
190            super::completion::reads_image(image.media_type.as_ref(), place) && carried(&image.data)
191        }
192        Media::Audio(audio) => audio.media_type.is_some() && carried(&audio.data),
193        Media::Video(video) => video.media_type.is_some() && carried(&video.data),
194        Media::Document(document) => match (&document.media_type, &document.data) {
195            (None, _) => false,
196            (Some(DocumentMediaType::PDF), data) => carried(data),
197            (Some(_), data) => matches!(data, Source::String(_)),
198        },
199    }
200}
201
202/// Read an existing interaction resource or resume its event stream.
203/// Unary mode retrieves the resource once. Streaming mode resumes after
204/// `last_event_id`, or from the beginning if no event id is supplied.
205#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
206pub struct InteractionResume {
207    /// The key and the API root.
208    pub provider: crate::providers::gemini::GeminiConfig,
209    /// The interaction to read.
210    pub interaction_id: String,
211    /// The last event the consumer saw, so a resumed stream does not
212    /// redeliver it. `None` resumes from the beginning, as the API defaults.
213    pub last_event_id: Option<String>,
214}
215
216impl InteractionResume {
217    /// The wire for the interaction `interaction_id`.
218    pub fn new(
219        provider: crate::providers::gemini::GeminiConfig,
220        interaction_id: impl Into<String>,
221    ) -> Self {
222        Self {
223            provider,
224            interaction_id: interaction_id.into(),
225            last_event_id: None,
226        }
227    }
228
229    /// Resume a streamed read after the event `last_event_id`.
230    pub fn after_event(mut self, last_event_id: impl Into<String>) -> Self {
231        self.last_event_id = Some(last_event_id.into());
232        self
233    }
234}
235
236impl crate::wire::Wire for InteractionResume {
237    type Op = crate::operation::Completion;
238    type Payload = crate::wire::Encoded;
239    type Frame = crate::wire::WireFrame;
240    type Decoder<'id> = streaming::InteractionsDecoder;
241    type Reassembler = streaming::document::Interaction;
242
243    /// The interaction names its own model; this wire addresses no model id.
244    /// The decoder reports the model the interaction names, which the turn's
245    /// origin takes.
246    fn describe(&self) -> Descriptor<'_> {
247        Descriptor::new(PROVIDER_NAME)
248            .telemetry(|_| GenAiOperation::Chat)
249            .replay(self)
250    }
251
252    /// Reads an existing interaction, so the request carries no body: what
253    /// to read is the wire's own data. A request that sets options,
254    /// provider options or `additional_params` is refused, since nothing
255    /// would send them.
256    fn encode(
257        &self,
258        request: CompletionRequest,
259        mode: Mode,
260    ) -> Result<crate::wire::Encoded, EncodeError> {
261        let params = crate::completion::options::request_params(
262            self,
263            &request,
264            // `request_params` takes `additional_params.tools` (provider
265            // tools included) out of the raw layer for the base to append,
266            // so the base refuses them here or they would vanish.
267            |input| {
268                if input.raw_tools()?.is_empty() {
269                    Ok(Map::new())
270                } else {
271                    Err(EncodeError::request(
272                        "a resumed interaction takes no `additional_params.tools` or provider \
273                         tools: it is read, not created",
274                    ))
275                }
276            },
277            crate::completion::options::RawAt::Top,
278            &[],
279        )?;
280        if !params.is_empty() {
281            return Err(EncodeError::request(
282                "a resumed interaction takes no provider options or `additional_params`: it is \
283                 read, not created",
284            ));
285        }
286        let id = &self.interaction_id;
287        let (path, framing) = match mode {
288            Mode::Unary => (
289                format!("/v1beta/interactions/{id}"),
290                crate::wire::Framing::Whole,
291            ),
292            Mode::Streaming => {
293                let mut query = form_urlencoded::Serializer::new(String::new());
294                query.append_pair("stream", "true");
295                if let Some(last_event_id) = &self.last_event_id {
296                    query.append_pair("last_event_id", last_event_id);
297                }
298                let path = format!("/v1beta/interactions/{id}?{}&alt=sse", query.finish());
299                (path, crate::wire::Framing::Sse)
300            }
301        };
302        let request = http::Request::get(self.provider.interactions_uri(&path))
303            .header(
304                crate::providers::gemini::GeminiConfig::INTERACTIONS_KEY_HEADER,
305                self.provider.api_key.expose(),
306            )
307            .body(crate::wire::Body::empty())?;
308        Ok(crate::wire::Encoded::new(request, framing))
309    }
310
311    fn decoder<'id>(&self) -> Self::Decoder<'id> {
312        streaming::InteractionsDecoder::default()
313    }
314}
315
316impl ReplayTarget for InteractionResume {
317    /// A resumed interaction is read, not created, so every set option is
318    /// refused.
319    fn map_options(
320        &self,
321        _request: &CompletionRequest,
322        fields: crate::completion::options::OptionFields<'_>,
323    ) -> crate::completion::options::OptionMap {
324        use crate::completion::options::{Mapping, OptionFields, OptionMap};
325        let OptionFields {
326            reasoning,
327            cache,
328            service_tier,
329            verbosity,
330            parallel_tool_calls,
331            top_p,
332            seed,
333            stop,
334        } = fields;
335        let refuse = |set: bool| match set {
336            true => Mapping::unsupported(
337                "a resumed interaction is read, not created; its options were fixed when it was \
338                 created",
339            ),
340            false => Mapping::Nothing,
341        };
342        OptionMap {
343            reasoning: refuse(reasoning.is_some()),
344            cache: refuse(cache.is_some()),
345            service_tier: refuse(service_tier.is_some()),
346            verbosity: refuse(verbosity.is_some()),
347            parallel_tool_calls: refuse(parallel_tool_calls.is_some()),
348            top_p: refuse(top_p.is_some()),
349            seed: refuse(seed.is_some()),
350            stop: refuse(!stop.is_empty()),
351        }
352    }
353
354    fn api(&self) -> crate::message::Api {
355        API
356    }
357
358    fn provider(&self) -> &str {
359        PROVIDER_NAME
360    }
361
362    fn model(&self) -> &str {
363        ""
364    }
365
366    fn accepts(&self, model: &str) -> crate::completion::Accepts {
367        super::completion::accepts(model)
368    }
369
370    fn normalize_tool_call_id(
371        &self,
372        id: &str,
373        _model: &str,
374        _: Option<&crate::message::Origin>,
375    ) -> String {
376        crate::providers::internal::wire_ids::legal_call_id(id, 64)
377    }
378
379    /// Gemini takes system text only in `systemInstruction`: later system
380    /// messages fold into the leading one, as pi's `collapseSystemMessages`.
381    fn later_system(&self, _model: &str) -> crate::completion::LaterSystem {
382        crate::completion::LaterSystem::Leading
383    }
384
385    fn call_id_slot(&self) -> Option<&'static str> {
386        Some("/id")
387    }
388}
389
390/// The create-interaction body `request` sends on `wire`: the wire's own
391/// encoding, then the mapped options, then `additional_params`, merged by
392/// [`request_params`](crate::completion::options::request_params). Its
393/// `generation_config` merges over the typed fields key by key, its `tools`
394/// add to the request's, and an `agent` takes the place of the model.
395/// System messages become the system instruction. `stream`, when set,
396/// overrides the one the parameters name.
397///
398/// # Errors
399///
400/// When an option is refused, `additional_params` is not an object, its
401/// `tools` are not an array, or it sets `response_format` without
402/// `response_mime_type`.
403pub(crate) fn create_request_body(
404    wire: &Interactions,
405    request: &CompletionRequest,
406    stream: Option<bool>,
407) -> Result<crate::completion::options::FinalBody, EncodeError> {
408    let rewrites: Vec<_> = stream
409        .map(crate::completion::options::Rewrite::Stream)
410        .into_iter()
411        .collect();
412    crate::completion::options::request_params(
413        wire,
414        request,
415        |input| create_base(wire, request, input),
416        crate::completion::options::RawAt::Top,
417        &rewrites,
418    )
419}
420
421/// The wire's own encoding of `request`: model, input steps, system
422/// instruction, tools and the typed fields in `generation_config`.
423fn create_base(
424    wire: &Interactions,
425    request: &CompletionRequest,
426    input: &mut crate::completion::options::BaseInput<'_>,
427) -> Result<Map<String, Value>, EncodeError> {
428    let model = request.model.clone().unwrap_or_else(|| wire.model.clone());
429    let set = |key: &str| input.param(key).is_some_and(|value| !value.is_null());
430    let agent = set("agent");
431    if set("response_format") && !set("response_mime_type") {
432        return Err(EncodeError::request(
433            "response_mime_type is required when response_format is set",
434        ));
435    }
436    let choice = request.tool_choice.clone().map(|choice| match choice {
437        Choice::Auto => json!("auto"),
438        Choice::None => json!("none"),
439        Choice::Required => json!("any"),
440        Choice::Specific { function_names } => {
441            json!({ "allowed_tools": { "mode": "validated", "tools": function_names } })
442        }
443    });
444    let config: Map<String, Value> = [
445        ("temperature", request.temperature.map(Value::from)),
446        ("max_output_tokens", request.max_tokens.map(Value::from)),
447        ("tool_choice", choice),
448    ]
449    .into_iter()
450    .filter_map(|(key, value)| Some((key.to_owned(), value?)))
451    .collect();
452    let mut tools: Vec<Value> = request
453        .tools
454        .iter()
455        .map(|tool| json!({ "type": "function", "name": tool.name, "description": tool.description, "parameters": tool.parameters }))
456        .collect();
457    tools.extend(input.raw_tools()?);
458    let (mut system, mut history) = (Vec::new(), Vec::new());
459    for message in &request.chat_history {
460        match message {
461            Message::System { content } => system.push(content.clone()),
462            message => history.push(message.clone()),
463        }
464    }
465    let mut body = Map::new();
466    if !config.is_empty() {
467        body.insert("generation_config".to_owned(), Value::Object(config));
468    }
469    if !tools.is_empty() {
470        body.insert("tools".to_owned(), Value::Array(tools));
471    }
472    if !system.is_empty() {
473        body.insert("system_instruction".to_owned(), json!(system.join("\n\n")));
474    }
475    if !agent {
476        body.insert("model".to_owned(), json!(model));
477    }
478    body.insert(
479        "input".to_owned(),
480        Value::Array(steps(history, wire, &model)?),
481    );
482    Ok(body)
483}
484
485/// The steps `history` sends: a user message's content grouped into
486/// `user_input` steps around its function results, each a step of its own,
487/// and one step per assistant block.
488fn steps(
489    history: Vec<Message>,
490    target: &dyn ReplayTarget,
491    model: &str,
492) -> Result<Vec<Value>, EncodeError> {
493    let ids = WireIds::for_target(&history, target, model);
494    let mut steps = Vec::new();
495    for message in history {
496        match message {
497            Message::System { content } => {
498                steps.push(json!({ "type": "user_input", "content": [{ "type": "text", "text": content }] }));
499            }
500            Message::User { content } => {
501                let mut run = Vec::new();
502                for part in content {
503                    let UserContent::ToolResult(result) = part else {
504                        run.push(user_content(part)?);
505                        continue;
506                    };
507                    if !run.is_empty() {
508                        steps.push(
509                            json!({ "type": "user_input", "content": std::mem::take(&mut run) }),
510                        );
511                    }
512                    let mut contents = result.content;
513                    let value = match (contents.len(), contents.pop()) {
514                        (1, Some(ToolResultContent::Text(text))) => Value::String(text.text),
515                        // A scalar or array JSON result is wrapped as the
516                        // generate wire wraps it: sent as a text block it is a
517                        // multimodal response, which the models refuse.
518                        (
519                            1,
520                            Some(ToolResultContent::Json {
521                                value: value @ (Value::String(_) | Value::Object(_)),
522                            }),
523                        ) => value,
524                        (1, Some(ToolResultContent::Json { value })) => json!({ "result": value }),
525                        (_, last) => {
526                            contents.extend(last);
527                            Value::Array(
528                                contents
529                                    .into_iter()
530                                    .map(result_content)
531                                    .collect::<Result<_, _>>()?,
532                            )
533                        }
534                    };
535                    let mut step = json!({
536                        "type": "function_result",
537                        "name": result.name,
538                        "call_id": ids.of(&result.call),
539                        "result": value,
540                    });
541                    if let (true, Some(step)) = (result.is_error, step.as_object_mut()) {
542                        step.insert("is_error".to_owned(), Value::Bool(true));
543                    }
544                    steps.push(step);
545                }
546                if !run.is_empty() {
547                    steps.push(json!({ "type": "user_input", "content": run }));
548                }
549            }
550            Message::Assistant(turn) => {
551                for block in &turn.content {
552                    steps.extend(assistant_step(block, target, &ids)?);
553                }
554            }
555        }
556    }
557    Ok(steps)
558}
559
560/// A block of a tool result as a content item.
561fn result_content(content: ToolResultContent) -> Result<Value, EncodeError> {
562    Ok(match content {
563        ToolResultContent::Text(text) => json!({ "type": "text", "text": text.text }),
564        ToolResultContent::Json { value } => json!({ "type": "text", "text": value.to_string() }),
565        ToolResultContent::Image(image) => media("image", image.media_type, image.data)?,
566    })
567}
568
569/// A user part as a content item. A text document goes as text, so that RAG
570/// context reads as prose.
571fn user_content(part: UserContent) -> Result<Value, EncodeError> {
572    match part {
573        UserContent::Text(text) => Ok(json!({ "type": "text", "text": text.text })),
574        UserContent::Image(image) => media("image", image.media_type, image.data),
575        UserContent::Audio(audio) => media("audio", audio.media_type, audio.data),
576        UserContent::Video(video) => media("video", video.media_type, video.data),
577        UserContent::Document(document) => match (document.media_type, document.data) {
578            (Some(media_type), Source::String(text)) if media_type != DocumentMediaType::PDF => {
579                Ok(json!({ "type": "text", "text": text }))
580            }
581            (media_type, data) => media("document", media_type, data),
582        },
583        UserContent::ToolResult(_) => {
584            Err(EncodeError::request("a tool result is a step of its own"))
585        }
586    }
587}
588
589/// A media content item of `kind`: a URL as its `uri`, and base64 data, or a
590/// string's bytes in base64, as its `data`. [`encodes`] refuses every other
591/// form, so the adapter passes none.
592fn media<M: MimeType>(
593    kind: &str,
594    media_type: Option<M>,
595    source: Source,
596) -> Result<Value, EncodeError> {
597    let unsendable = || {
598        EncodeError::request(format!(
599            "Gemini Interactions cannot receive this {kind} in its form"
600        ))
601    };
602    let mime_type = media_type.ok_or_else(unsendable)?.to_mime_type().to_owned();
603    let key = match super::completion::carried(source, false)? {
604        (true, uri) => ("uri", uri),
605        (false, data) => ("data", data),
606    };
607    Ok(json!({ "type": kind, key.0: key.1, "mime_type": mime_type }))
608}
609
610/// An assistant block as its step: the provider's step while the block is
611/// current, else one rebuilt from its canonical fields, keeping the keys
612/// its edited step names; `None` for reasoning with nothing to send.
613fn assistant_step(
614    block: &AssistantContent,
615    target: &dyn ReplayTarget,
616    ids: &WireIds,
617) -> Result<Option<Value>, EncodeError> {
618    let identity = match block.replay(target, ids) {
619        Replay::Item(item) => return Ok(Some(item.into_owned())),
620        Replay::Identity(identity) => identity,
621        Replay::Rebuild => Map::new(),
622    };
623    let output = |content: Value| json!({ "type": "model_output", "content": [content] });
624    let mut step = match block {
625        AssistantContent::Text(text) => output(json!({ "type": "text", "text": text.text })),
626        AssistantContent::Reasoning(reasoning)
627            if reasoning.redacted || reasoning.text.trim().is_empty() =>
628        {
629            return Ok(None);
630        }
631        AssistantContent::Reasoning(reasoning) => {
632            json!({ "type": "thought", "summary": [{ "type": "text", "text": reasoning.text }] })
633        }
634        AssistantContent::ToolCall(call) => json!({
635            "type": "function_call",
636            "name": call.function.name,
637            "arguments": call.function.arguments_value(),
638        }),
639        AssistantContent::Image(image) => output(media(
640            "image",
641            image.media_type.clone(),
642            image.data.clone(),
643        )?),
644        AssistantContent::Opaque(opaque) => return Ok(Some(opaque.item.clone())),
645    };
646    if let Some(fields) = step.as_object_mut() {
647        fields.extend(identity);
648        if let AssistantContent::ToolCall(call) = block {
649            fields.insert("id".to_owned(), json!(ids.of(&call.id)));
650        }
651    }
652    Ok(Some(step))
653}
654
655#[cfg(test)]
656mod tests;
657
658#[cfg(test)]
659mod history_tests;