Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,12 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [Unreleased]

### Changed

- [**breaking**] Rig-owned public data-bearing Serde enums use explicit discriminators; multi-turn stream events use nested adjacent tags, and intentional unit enums use explicit lowercase strings. Old domain/event representations do not load; provider HTTP/SSE wire JSON is unchanged except for the Cohere `tool_choice` correction below. See the crate changelogs and MIGRATING for exact shapes.
- *(cohere)* [**fix**] translate normalized tool choices to Cohere's native v2 Chat schema: omit `tool_choice` for auto mode, send `"NONE"`/`"REQUIRED"` scalar values, and implement named choices by filtering the advertised tools and requiring a call.

## [0.41.0](https://github.com/0xPlaygrounds/rig/compare/v0.40.0...v0.41.0) - 2026-07-28

### Added
Expand Down
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

116 changes: 99 additions & 17 deletions MIGRATING.md
Original file line number Diff line number Diff line change
Expand Up @@ -1212,9 +1212,11 @@ setters rather than at provider call sites.

Both invariants also hold through `Deserialize`: the two types deserialize
via a wire-shape mirror that funnels through `new(...)` and the setters, so
a persisted `"finish_reason": "stop"` alongside a tool-call choice comes
back as `ToolCalls` and a persisted `""` identifier comes back as `None`.
The serialized wire format is unchanged.
a persisted `"finish_reason": {"type": "stop"}` alongside a tool-call choice
comes back as `ToolCalls` and a persisted `""` identifier comes back as
`None`. The normalized finish-reason representation itself changes in this
release; see "Rig-owned serializable enums have explicit discriminators"
below. Provider-native finish-reason fields remain unchanged.

Corrupt stream frames (payloads that are not valid JSON) are now surfaced as
`Err` items on the stream instead of being logged and silently skipped; the
Expand Down Expand Up @@ -1354,20 +1356,100 @@ Run it once over your store at migration time; the runtime load path stays
tolerant on purpose (loudness at the migration boundary, not on every
restore).

**Streaming events**: `StreamedAssistantContent` is a tolerant decode with an
`Unknown` catch-all. A stream item whose text block carries stray sibling
keys — 0.41's flatten shape, or a relay stamping bookkeeping keys onto text
items — decodes as stream *text* with the stray keys dropped: the text is
assembled, nothing is excluded. A replayed **tagged** assistant block
(`toolcall`/`reasoning`/`image` — the tagged `AssistantContent`
serialization is not a stream-item shape) decodes as `Unknown` and is
excluded from assembly, as is a text item whose `additional_params` is
malformed (a non-object — the strict decode rejects the known field); the
**agent assembler** — the one rig component that ingests replayed stream
events — counts both kinds of exclusion and logs a single `tracing` warning
per turn, on every termination path
(`StreamedTurnAssembler::excluded_assistant_content` exposes the count). A consumer assembling self-deserialized events with its
own logic gets no warning and should apply the same check itself.
**Streaming events** now declare their identity instead of inferring it from
object shape. `StreamedAssistantContent` and `StreamedUserContent` use adjacent
snake-case `type`/`content` fields, while `MultiTurnStreamItem` uses the same
adjacent layout with camel-case outer tags:

```json
{
"type": "streamAssistantItem",
"content": {
"type": "text",
"content": {"text": "hello"}
}
}
```

An unmodeled provider event is representable only through the explicit
`Unknown` variant. Its raw value is preserved inside the adjacent wrapper,
even when the payload has keys that resemble a known Rig event:

```json
{
"type": "unknown",
"content": {
"type": "provider.future_event",
"text": "opaque provider content",
"kind": "final"
}
}
```

Provider adapters classify native frames first and construct this variant;
arbitrary provider JSON is not deserialized directly into a Rig stream event.
A known Rig tag with malformed content is a decode error and never falls
through to `Unknown`.

### Rig-owned serializable enums have explicit discriminators

The remaining public, data-bearing Rig enums now use explicit tags too. These
are clean serialization breaks: every old form in the table fails to load, and
the current form is the only accepted representation.

| Type | Before | Current |
|---|---|---|
| `ToolCallDeltaContent` | `{"Name":"lookup"}` | `{"type":"name","content":"lookup"}` |
| `StreamedAssistantContent` | `{"text":"hello"}` | `{"type":"text","content":{"text":"hello"}}` |
| `StreamedUserContent` | `{"tool_result":{"call":"call_1","name":"lookup","content":[{"type":"text","text":"ok"}]},"internal_call_id":"internal_1"}` | `{"type":"tool_result","content":{"tool_result":{"call":"call_1","name":"lookup","content":[{"type":"text","text":"ok"}]},"internal_call_id":"internal_1"}}` |
| `MultiTurnStreamItem` | `{"type":"streamAssistantItem","text":"hello"}` | `{"type":"streamAssistantItem","content":{"type":"text","content":{"text":"hello"}}}` |
| `MediaType` | `{"Image":"png"}` | `{"type":"image","content":"png"}` |
| `ToolChoice` | `"auto"`; `{"specific":{"function_names":["lookup"]}}` | `{"type":"auto"}`; `{"type":"specific","function_names":["lookup"]}` |
| `FinishReason` | `"stop"`; `{"other":"RECITATION"}` | `{"type":"stop"}`; `{"type":"other","content":"RECITATION"}` |
| `Filter<V>` | `{"eq":["status","ready"]}` | `{"type":"eq","content":["status","ready"]}` |
| `ModelListingError` | `{"ApiError":{"status_code":429,"message":"limited"}}` | `{"type":"api_error","status_code":429,"message":"limited"}` |
| `OutputMode` | `"Auto"` | `"auto"` |
| `SqliteDistanceMetric` | `"Cosine"` | `"cosine"` |

The adjacent forms use `content` for newtype and tuple payloads; unit variants
omit it. **Struct variants also nest under `content`** — that covers most of the
retagged stream enums (`StreamedAssistantContent::{ToolCall, ToolCallDelta,
Reasoning, ReasoningDelta}` and `MultiTurnStreamItem::{ToolExecutionCommitted,
ModelTurnRetried}`), whose fields move inside `content` rather than sitting
beside `type`. `ToolChoice` and `ModelListingError` are the exception: they use
an internal `type` field, so their struct-variant fields stay at the top level
(see the `ToolChoice` row above).
`OutputMode` and `SqliteDistanceMetric` remain intentional string-valued unit
enums, but their spellings are now explicit rather than Rust's PascalCase
names.

Persistence impact follows where these values are embedded:

- Persisted `AgentRun` and directly serialized `CompletionRequest` values carry
`ToolChoice` and require its tagged form.
- Serialized `CompletionResponse` and `StreamFinal` values carry the new
`FinishReason` form. `PromptResponse` has no new field shape from this
change, but a `MultiTurnStreamItem::FinalResponse` stream-log record now
wraps the whole response under the outer `content` field.
- The public OpenAI-compatible raw terminal record
`openai::completion::streaming::StreamingCompletionResponse<U>` carries the
normalized `FinishReason` form too. This also affects its Mistral, Groq,
DeepSeek, and OpenRouter aliases when those values are serialized directly.
- Recorded multi-turn stream events, standalone media-type values, serialized
vector-store filters, model-listing errors, agent output-mode values, and
SQLite metric configuration need the corresponding rewrite above.

Provider HTTP request/response and SSE payload JSON is otherwise unchanged.
Provider-native wire enums retain the upstream schema (including genuine
untagged unions and scalar finish reasons); conversions map between those wire
types and these Rig-owned values at the boundary. The exception is a Cohere v2
Chat correction:
`ToolChoice::Auto` now omits `tool_choice`, `None` and `Required` serialize as the
native `"NONE"` and `"REQUIRED"` scalars, and `Specific` filters the advertised
tools to the requested names and sends `"REQUIRED"`. Existing explicitly tagged
Rig domain enums (`Message`, `UserContent`, `AssistantContent`,
`ReasoningContent`, `ToolResultContent`, and `DocumentSourceKind`) are
unchanged.

### Two pre-`Vec` serde accommodations are gone

Expand Down
4 changes: 3 additions & 1 deletion crates/rig-agent/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Changed

- *(agent)* [**breaking**] `MultiTurnStreamItem` now uses adjacent camel-case `"type"`/`"content"` tagging around the newly tagged `StreamedAssistantContent` and `StreamedUserContent` payloads, eliminating nested tag collisions and structural inference. Old internal/untagged stream-log records do not load. Persisted `AgentRun::tool_choice` uses rig-core's tagged `ToolChoice` object, and `OutputMode` now serializes as the explicit lowercase strings `"auto"`, `"tool"`, `"native"`, and `"prompted"` instead of PascalCase. `PromptResponse` itself gains no field-shape change here, but a `MultiTurnStreamItem::FinalResponse` record now wraps it under the outer `content` field; see MIGRATING.

- *(agent)* [**breaking**] persisted histories and `AgentRun`/`PromptResponse` JSON carry rig-core's tagged assistant content (`{"type": "text", ...}`); the untagged shape does not load — see rig-core's entry and MIGRATING. The flatten `Some({})` round-trip artifact is gone, so `is_empty_assistant_turn`'s classification is identical before and after a persist/restore with no special-casing

- *(agent)* [**behavior**] the streamed assembler counts the stream items it excludes from assembly that carry assistant content — replayed tagged assistant blocks (the tagged `AssistantContent` serialization is not a stream-item shape) and text items whose `additional_params` is malformed — and logs a single warning per turn, on every termination path, instead of one per stream item; `StreamedTurnAssembler::excluded_assistant_content` exposes the count, and the full decode-outcome contract is pinned by an enum-driven matrix test. A stream item whose text block carries stray sibling keys decodes as stream *text* — the text is assembled and only the stray keys drop
- *(agent)* [**behavior**] the streamed assembler counts explicitly constructed `StreamedAssistantContent::Unknown` items that carry Rig-like assistant content and logs one warning per turn, on every termination path; `StreamedTurnAssembler::excluded_assistant_content` exposes the count. Public stream-event deserialization is now strict and tagged, so malformed known items fail at decode and provider-native items reach the assembler as `Unknown` only after provider classification.

- *(agent)* [**breaking**] `PromptResponse::content` returns `&[AssistantContent]` instead of `&Vec<AssistantContent>`, matching its slice-returning siblings

Expand Down
173 changes: 148 additions & 25 deletions crates/rig-agent/src/agent/prompt_request/streaming.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,8 +46,14 @@ pub type StreamingResult =
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub type StreamingResult = Pin<Box<dyn Stream<Item = Result<MultiTurnStreamItem, StreamingError>>>>;

/// One event emitted by a multi-turn agent stream.
///
/// The public JSON format is adjacently tagged with camel-case variant names.
/// Nested assistant and user events carry their own adjacent discriminator, so
/// an assistant text item is represented as
/// `{"type":"streamAssistantItem","content":{"type":"text","content":{"text":"hello"}}}`.
#[derive(Deserialize, Serialize, Debug, Clone)]
#[serde(tag = "type", rename_all = "camelCase")]
#[serde(tag = "type", content = "content", rename_all = "camelCase")]
#[non_exhaustive]
pub enum MultiTurnStreamItem {
/// A streamed assistant content item — the content the **model emitted**:
Expand Down Expand Up @@ -1643,6 +1649,117 @@ pub async fn stream_to_stdout(
Ok(final_res)
}

#[cfg(test)]
mod stream_item_serde_tests {
use super::{CompletionCall, MultiTurnStreamItem, PromptResponse};
use crate::streaming::{StreamedAssistantContent, StreamedUserContent};
use rig_core::{
completion::Usage,
message::{ToolCall, ToolCallId, ToolFunction, ToolResult, ToolResultContent},
};
use serde_json::{Value, json};

fn tool_call() -> ToolCall {
ToolCall::new(
ToolCallId::new_or_mint("call_1"),
ToolFunction::new("lookup".to_owned(), json!({"query": "rig"})),
)
}

fn tool_result() -> ToolResult {
ToolResult {
call: ToolCallId::new_or_mint("call_1"),
provider: None,
name: "lookup".to_owned(),
content: vec![ToolResultContent::text("found")],
}
}

fn assert_wire_round_trip(item: MultiTurnStreamItem, expected: Value) {
let encoded = serde_json::to_value(&item).expect("serialize multi-turn stream item");
assert_eq!(encoded, expected);
let decoded: MultiTurnStreamItem =
serde_json::from_value(encoded).expect("deserialize multi-turn stream item");
assert_eq!(
serde_json::to_value(decoded).expect("re-serialize multi-turn stream item"),
expected
);
}

#[test]
fn every_multi_turn_stream_variant_uses_an_adjacent_payload() {
let assistant = StreamedAssistantContent::text("hello");
let assistant_json = serde_json::to_value(&assistant).expect("serialize assistant item");
assert_wire_round_trip(
MultiTurnStreamItem::StreamAssistantItem(assistant),
json!({
"type": "streamAssistantItem",
"content": assistant_json
}),
);

let call = tool_call();
assert_wire_round_trip(
MultiTurnStreamItem::ToolExecutionCommitted {
tool_call: call.clone(),
internal_call_id: "internal_1".to_owned(),
},
json!({
"type": "toolExecutionCommitted",
"content": {
"tool_call": serde_json::to_value(&call).expect("serialize tool call"),
"internal_call_id": "internal_1"
}
}),
);

let result = tool_result();
let user = StreamedUserContent::ToolResult {
tool_result: result,
internal_call_id: "internal_1".to_owned(),
};
let user_json = serde_json::to_value(&user).expect("serialize user item");
assert_wire_round_trip(
MultiTurnStreamItem::StreamUserItem(user),
json!({"type": "streamUserItem", "content": user_json}),
);

let completion_call = CompletionCall::new(2, Usage::new());
let completion_call_json =
serde_json::to_value(completion_call).expect("serialize completion call");
assert_wire_round_trip(
MultiTurnStreamItem::CompletionCall(completion_call),
json!({"type": "completionCall", "content": completion_call_json}),
);

assert_wire_round_trip(
MultiTurnStreamItem::ModelTurnRetried { turn: 3 },
json!({"type": "modelTurnRetried", "content": {"turn": 3}}),
);

let response = PromptResponse::new("done", Usage::new());
let response_json = serde_json::to_value(&response).expect("serialize prompt response");
assert_wire_round_trip(
MultiTurnStreamItem::FinalResponse(response),
json!({"type": "finalResponse", "content": response_json}),
);
}

#[test]
fn old_internal_stream_shapes_and_malformed_known_tags_fail() {
for old_or_invalid in [
json!({"type": "streamAssistantItem", "text": "old"}),
json!({"type": "modelTurnRetried", "turn": 3}),
json!({"type": "completionCall", "call_index": 3, "usage": null}),
json!({"content": {"turn": 3}}),
json!({"type": "future", "content": {}}),
json!({"type": "modelTurnRetried", "content": {"turn": "three"}}),
] {
assert!(serde_json::from_value::<MultiTurnStreamItem>(old_or_invalid).is_err());
}
}
}

#[cfg(test)]
#[allow(irrefutable_let_patterns, unreachable_patterns)]
mod migrated_tests {
Expand Down Expand Up @@ -3067,15 +3184,17 @@ mod migrated_tests {
value,
serde_json::json!({
"type": "completionCall",
"call_index": 2,
"usage": {
"input_tokens": 3,
"output_tokens": 4,
"total_tokens": 7,
"cached_input_tokens": 0,
"cache_creation_input_tokens": 0,
"tool_use_prompt_tokens": 0,
"reasoning_tokens": 0,
"content": {
"call_index": 2,
"usage": {
"input_tokens": 3,
"output_tokens": 4,
"total_tokens": 7,
"cached_input_tokens": 0,
"cache_creation_input_tokens": 0,
"tool_use_prompt_tokens": 0,
"reasoning_tokens": 0,
}
}
})
);
Expand All @@ -3099,27 +3218,31 @@ mod migrated_tests {
value,
serde_json::json!({
"type": "completionCall",
"call_index": 3,
"usage": {
"input_tokens": 0,
"output_tokens": 0,
"total_tokens": 0,
"cached_input_tokens": 0,
"cache_creation_input_tokens": 0,
"tool_use_prompt_tokens": 0,
"reasoning_tokens": 0,
"content": {
"call_index": 3,
"usage": {
"input_tokens": 0,
"output_tokens": 0,
"total_tokens": 0,
"cached_input_tokens": 0,
"cache_creation_input_tokens": 0,
"tool_use_prompt_tokens": 0,
"reasoning_tokens": 0,
}
}
})
);

// Stream items serialized before the Option encoding was dropped used
// `"usage": null`; they must still deserialize.
// CompletionCall's null-usage tolerance still applies inside the new
// adjacent event payload.
let legacy: MultiTurnStreamItem = serde_json::from_value(serde_json::json!({
"type": "completionCall",
"call_index": 3,
"usage": null
"content": {
"call_index": 3,
"usage": null
}
}))
.expect("legacy null-usage event should deserialize");
.expect("null-usage event should deserialize");
match legacy {
MultiTurnStreamItem::CompletionCall(call) => {
assert_eq!(call, CompletionCall::new(3, Usage::new()));
Expand Down Expand Up @@ -3147,7 +3270,7 @@ mod migrated_tests {
let value = serde_json::to_value(&item).expect("serialize final response");

assert_eq!(
value.get("completion_calls"),
value.pointer("/content/completion_calls"),
Some(&serde_json::json!([
{
"call_index": 0,
Expand Down
Loading
Loading