use std::collections::HashSet;
use axum::response::sse::Event;
use dynamo_protocols::types::ChatCompletionMessageContent;
use uuid::Uuid;
use super::types::{
AnthropicDelta, AnthropicErrorBody, AnthropicMessageDeltaBody, AnthropicMessageResponse,
AnthropicResponseContentBlock, AnthropicStopReason, AnthropicStreamEvent, AnthropicUsage,
};
use crate::protocols::openai::chat_completions::NvCreateChatCompletionStreamResponse;
use crate::protocols::unified::AnthropicContext;
pub struct AnthropicStreamConverter {
model: String,
message_id: String,
api_context: Option<AnthropicContext>,
thinking_block_started: bool,
thinking_block_closed: bool,
thinking_block_index: u32,
text_block_started: bool,
text_block_closed: bool,
text_block_index: u32,
input_token_count: u32,
output_token_count: u32,
cached_token_count: Option<u32>,
tool_call_states: Vec<ToolCallState>,
tool_calls_sent: HashSet<String>,
next_block_index: u32,
stop_reason: Option<AnthropicStopReason>,
}
struct ToolCallState {
id: String,
name: String,
accumulated_args: String,
block_index: u32,
started: bool,
stopped: bool,
}
impl AnthropicStreamConverter {
pub fn new(model: String, estimated_input_tokens: u32) -> Self {
Self {
model,
message_id: format!("msg_{}", Uuid::new_v4().simple()),
api_context: None,
thinking_block_started: false,
thinking_block_closed: false,
thinking_block_index: 0,
text_block_started: false,
text_block_closed: false,
text_block_index: 0,
input_token_count: estimated_input_tokens,
output_token_count: 0,
cached_token_count: None,
tool_call_states: Vec::new(),
tool_calls_sent: HashSet::new(),
next_block_index: 0,
stop_reason: None,
}
}
pub fn with_context(
model: String,
estimated_input_tokens: u32,
context: AnthropicContext,
) -> Self {
let mut converter = Self::new(model, estimated_input_tokens);
converter.api_context = Some(context);
converter
}
pub fn emit_start_events(&mut self) -> Vec<Result<Event, anyhow::Error>> {
let mut events = Vec::with_capacity(1);
self.append_start_events(&mut events);
events
}
pub fn append_start_events(&mut self, events: &mut Vec<Result<Event, anyhow::Error>>) {
let message = AnthropicMessageResponse {
id: self.message_id.clone(),
object_type: "message".to_string(),
role: "assistant".to_string(),
content: vec![],
model: self.model.clone(),
stop_reason: None,
stop_sequence: None,
usage: AnthropicUsage {
input_tokens: self.input_token_count,
output_tokens: 0,
cache_creation_input_tokens: None,
cache_read_input_tokens: None,
},
};
let event = AnthropicStreamEvent::MessageStart { message };
events.push(make_sse_event("message_start", &event));
}
pub fn process_chunk(
&mut self,
chunk: &NvCreateChatCompletionStreamResponse,
) -> Vec<Result<Event, anyhow::Error>> {
let mut events = Vec::new();
self.append_chunk_events(chunk, &mut events);
events
}
pub fn append_chunk_events(
&mut self,
chunk: &NvCreateChatCompletionStreamResponse,
events: &mut Vec<Result<Event, anyhow::Error>>,
) {
if let Some(usage) = &chunk.inner.usage {
self.output_token_count = usage.completion_tokens;
self.cached_token_count = usage
.prompt_tokens_details
.as_ref()
.and_then(|d| d.cached_tokens);
}
for choice in &chunk.inner.choices {
let delta = &choice.delta;
if let Some(ref fr) = choice.finish_reason {
self.stop_reason = Some(match fr {
dynamo_protocols::types::FinishReason::Stop => AnthropicStopReason::EndTurn,
dynamo_protocols::types::FinishReason::Length => AnthropicStopReason::MaxTokens,
dynamo_protocols::types::FinishReason::ToolCalls => {
AnthropicStopReason::ToolUse
}
dynamo_protocols::types::FinishReason::ContentFilter => {
AnthropicStopReason::EndTurn
}
dynamo_protocols::types::FinishReason::FunctionCall => {
AnthropicStopReason::ToolUse
}
});
}
if let Some(ref reasoning) = delta.reasoning_content
&& !reasoning.is_empty()
{
if !self.thinking_block_started {
self.thinking_block_started = true;
self.thinking_block_index = self.next_block_index;
self.next_block_index += 1;
let block_start = AnthropicStreamEvent::ContentBlockStart {
index: self.thinking_block_index,
content_block: AnthropicResponseContentBlock::Thinking {
thinking: String::new(),
signature: String::new(),
},
};
events.push(make_sse_event("content_block_start", &block_start));
}
let block_delta = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::ThinkingDelta {
thinking: reasoning.clone(),
},
};
events.push(make_sse_event("content_block_delta", &block_delta));
}
let content_text = match &delta.content {
Some(ChatCompletionMessageContent::Text(text)) => Some(text.as_str()),
_ => None,
};
if let Some(text) = content_text
&& !text.is_empty()
{
if self.thinking_block_started && !self.thinking_block_closed {
self.thinking_block_closed = true;
let sig_delta = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::SignatureDelta {
signature: "erased".to_string(),
},
};
events.push(make_sse_event("content_block_delta", &sig_delta));
let block_stop = AnthropicStreamEvent::ContentBlockStop {
index: self.thinking_block_index,
};
events.push(make_sse_event("content_block_stop", &block_stop));
}
if !self.text_block_started {
self.text_block_started = true;
self.text_block_index = self.next_block_index;
self.next_block_index += 1;
let block_start = AnthropicStreamEvent::ContentBlockStart {
index: self.text_block_index,
content_block: AnthropicResponseContentBlock::Text {
text: String::new(),
citations: None,
},
};
events.push(make_sse_event("content_block_start", &block_start));
}
let block_delta = AnthropicStreamEvent::ContentBlockDelta {
index: self.text_block_index,
delta: AnthropicDelta::TextDelta {
text: text.to_string(),
},
};
events.push(make_sse_event("content_block_delta", &block_delta));
}
if let Some(tool_calls) = &delta.tool_calls {
if self.thinking_block_started && !self.thinking_block_closed {
self.thinking_block_closed = true;
let sig_delta = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::SignatureDelta {
signature: "erased".to_string(),
},
};
events.push(make_sse_event("content_block_delta", &sig_delta));
let block_stop = AnthropicStreamEvent::ContentBlockStop {
index: self.thinking_block_index,
};
events.push(make_sse_event("content_block_stop", &block_stop));
}
if self.text_block_started && !self.text_block_closed {
self.text_block_closed = true;
let block_stop = AnthropicStreamEvent::ContentBlockStop {
index: self.text_block_index,
};
events.push(make_sse_event("content_block_stop", &block_stop));
}
for tc in tool_calls {
let tc_index = tc.index as usize;
while self.tool_call_states.len() <= tc_index {
let block_index = self.next_block_index;
self.next_block_index += 1;
self.tool_call_states.push(ToolCallState {
id: String::new(),
name: String::new(),
accumulated_args: String::new(),
block_index,
started: false,
stopped: false,
});
}
if let Some(id) = &tc.id {
self.tool_call_states[tc_index].id = id.clone();
}
if let Some(func) = &tc.function {
if let Some(name) = &func.name {
self.tool_call_states[tc_index].name = name.clone();
}
if let Some(args) = &func.arguments {
if !self.tool_call_states[tc_index].started {
let tc_id = self.tool_call_states[tc_index].id.clone();
if !tc_id.is_empty() && self.tool_calls_sent.contains(&tc_id) {
continue;
}
self.tool_call_states[tc_index].started = true;
let block_index = self.tool_call_states[tc_index].block_index;
let tc_name = self.tool_call_states[tc_index].name.clone();
if !tc_id.is_empty() {
self.tool_calls_sent.insert(tc_id.clone());
}
let block_start = AnthropicStreamEvent::ContentBlockStart {
index: block_index,
content_block: AnthropicResponseContentBlock::ToolUse {
id: tc_id,
name: tc_name,
input: serde_json::json!({}),
},
};
events.push(make_sse_event("content_block_start", &block_start));
}
self.tool_call_states[tc_index]
.accumulated_args
.push_str(args);
let block_index = self.tool_call_states[tc_index].block_index;
let block_delta = AnthropicStreamEvent::ContentBlockDelta {
index: block_index,
delta: AnthropicDelta::InputJsonDelta {
partial_json: args.clone(),
},
};
events.push(make_sse_event("content_block_delta", &block_delta));
if tc.id.is_some()
&& func.name.is_some()
&& !self.tool_call_states[tc_index].stopped
{
self.tool_call_states[tc_index].stopped = true;
let block_stop =
AnthropicStreamEvent::ContentBlockStop { index: block_index };
events.push(make_sse_event("content_block_stop", &block_stop));
}
}
}
}
}
}
}
pub fn emit_end_events(&mut self) -> Vec<Result<Event, anyhow::Error>> {
let mut events = Vec::new();
self.append_end_events(&mut events);
events
}
pub fn append_end_events(&mut self, events: &mut Vec<Result<Event, anyhow::Error>>) {
if self.thinking_block_started && !self.thinking_block_closed {
self.thinking_block_closed = true;
let sig_delta = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::SignatureDelta {
signature: "erased".to_string(),
},
};
events.push(make_sse_event("content_block_delta", &sig_delta));
let block_stop = AnthropicStreamEvent::ContentBlockStop {
index: self.thinking_block_index,
};
events.push(make_sse_event("content_block_stop", &block_stop));
}
if self.text_block_started && !self.text_block_closed {
let block_stop = AnthropicStreamEvent::ContentBlockStop {
index: self.text_block_index,
};
events.push(make_sse_event("content_block_stop", &block_stop));
}
for tc in &self.tool_call_states {
if tc.started && !tc.stopped {
let block_stop = AnthropicStreamEvent::ContentBlockStop {
index: tc.block_index,
};
events.push(make_sse_event("content_block_stop", &block_stop));
}
}
let message_delta = AnthropicStreamEvent::MessageDelta {
delta: AnthropicMessageDeltaBody {
stop_reason: self.stop_reason.clone(),
stop_sequence: None,
},
usage: AnthropicUsage {
input_tokens: self.input_token_count,
output_tokens: self.output_token_count,
cache_creation_input_tokens: None,
cache_read_input_tokens: self.cached_token_count,
},
};
events.push(make_sse_event("message_delta", &message_delta));
let message_stop = AnthropicStreamEvent::MessageStop {};
events.push(make_sse_event("message_stop", &message_stop));
}
pub fn emit_error_events(&mut self) -> Vec<Result<Event, anyhow::Error>> {
let mut events = Vec::with_capacity(1);
self.append_error_events(&mut events);
events
}
pub fn append_error_events(&mut self, events: &mut Vec<Result<Event, anyhow::Error>>) {
let error_event = AnthropicStreamEvent::Error {
error: AnthropicErrorBody {
error_type: "api_error".to_string(),
message: "An internal error occurred during generation.".to_string(),
},
};
events.push(make_sse_event("error", &error_event));
}
}
fn make_sse_event(event_type: &str, event: &AnthropicStreamEvent) -> Result<Event, anyhow::Error> {
let data = serde_json::to_string(event)?;
Ok(Event::default().event(event_type).data(data))
}
#[cfg(test)]
#[derive(Debug)]
struct TaggedEvent {
event_type: String,
data: AnthropicStreamEvent,
}
#[cfg(test)]
fn make_tagged_event(event_type: &str, event: &AnthropicStreamEvent) -> TaggedEvent {
TaggedEvent {
event_type: event_type.to_string(),
data: event.clone(),
}
}
#[cfg(test)]
impl AnthropicStreamConverter {
fn process_chunk_tagged(
&mut self,
chunk: &NvCreateChatCompletionStreamResponse,
) -> Vec<TaggedEvent> {
let mut events = Vec::new();
if let Some(usage) = &chunk.inner.usage {
self.output_token_count = usage.completion_tokens;
self.cached_token_count = usage
.prompt_tokens_details
.as_ref()
.and_then(|d| d.cached_tokens);
}
for choice in &chunk.inner.choices {
let delta = &choice.delta;
if let Some(ref fr) = choice.finish_reason {
self.stop_reason = Some(match fr {
dynamo_protocols::types::FinishReason::Stop => AnthropicStopReason::EndTurn,
dynamo_protocols::types::FinishReason::Length => AnthropicStopReason::MaxTokens,
dynamo_protocols::types::FinishReason::ToolCalls => {
AnthropicStopReason::ToolUse
}
dynamo_protocols::types::FinishReason::ContentFilter => {
AnthropicStopReason::EndTurn
}
dynamo_protocols::types::FinishReason::FunctionCall => {
AnthropicStopReason::ToolUse
}
});
}
if let Some(ref reasoning) = delta.reasoning_content
&& !reasoning.is_empty()
{
if !self.thinking_block_started {
self.thinking_block_started = true;
self.thinking_block_index = self.next_block_index;
self.next_block_index += 1;
let ev = AnthropicStreamEvent::ContentBlockStart {
index: self.thinking_block_index,
content_block: AnthropicResponseContentBlock::Thinking {
thinking: String::new(),
signature: String::new(),
},
};
events.push(make_tagged_event("content_block_start", &ev));
}
let ev = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::ThinkingDelta {
thinking: reasoning.clone(),
},
};
events.push(make_tagged_event("content_block_delta", &ev));
}
let content_text = match &delta.content {
Some(ChatCompletionMessageContent::Text(text)) => Some(text.as_str()),
_ => None,
};
if let Some(text) = content_text
&& !text.is_empty()
{
if self.thinking_block_started && !self.thinking_block_closed {
self.thinking_block_closed = true;
let ev = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::SignatureDelta {
signature: "erased".to_string(),
},
};
events.push(make_tagged_event("content_block_delta", &ev));
let ev = AnthropicStreamEvent::ContentBlockStop {
index: self.thinking_block_index,
};
events.push(make_tagged_event("content_block_stop", &ev));
}
if !self.text_block_started {
self.text_block_started = true;
self.text_block_index = self.next_block_index;
self.next_block_index += 1;
let ev = AnthropicStreamEvent::ContentBlockStart {
index: self.text_block_index,
content_block: AnthropicResponseContentBlock::Text {
text: String::new(),
citations: None,
},
};
events.push(make_tagged_event("content_block_start", &ev));
}
self.output_token_count += 1;
let ev = AnthropicStreamEvent::ContentBlockDelta {
index: self.text_block_index,
delta: AnthropicDelta::TextDelta {
text: text.to_string(),
},
};
events.push(make_tagged_event("content_block_delta", &ev));
}
if let Some(tool_calls) = &delta.tool_calls {
if self.thinking_block_started && !self.thinking_block_closed {
self.thinking_block_closed = true;
let ev = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::SignatureDelta {
signature: "erased".to_string(),
},
};
events.push(make_tagged_event("content_block_delta", &ev));
let ev = AnthropicStreamEvent::ContentBlockStop {
index: self.thinking_block_index,
};
events.push(make_tagged_event("content_block_stop", &ev));
}
if self.text_block_started && !self.text_block_closed {
self.text_block_closed = true;
let ev = AnthropicStreamEvent::ContentBlockStop {
index: self.text_block_index,
};
events.push(make_tagged_event("content_block_stop", &ev));
}
for tc in tool_calls {
let tc_index = tc.index as usize;
while self.tool_call_states.len() <= tc_index {
let block_index = self.next_block_index;
self.next_block_index += 1;
self.tool_call_states.push(ToolCallState {
id: String::new(),
name: String::new(),
accumulated_args: String::new(),
block_index,
started: false,
stopped: false,
});
}
if let Some(id) = &tc.id {
self.tool_call_states[tc_index].id = id.clone();
}
if let Some(func) = &tc.function {
if let Some(name) = &func.name {
self.tool_call_states[tc_index].name = name.clone();
}
if let Some(args) = &func.arguments {
if !self.tool_call_states[tc_index].started {
let tc_id = self.tool_call_states[tc_index].id.clone();
if !tc_id.is_empty() && self.tool_calls_sent.contains(&tc_id) {
continue;
}
self.tool_call_states[tc_index].started = true;
let block_index = self.tool_call_states[tc_index].block_index;
let tc_name = self.tool_call_states[tc_index].name.clone();
if !tc_id.is_empty() {
self.tool_calls_sent.insert(tc_id.clone());
}
let ev = AnthropicStreamEvent::ContentBlockStart {
index: block_index,
content_block: AnthropicResponseContentBlock::ToolUse {
id: tc_id,
name: tc_name,
input: serde_json::json!({}),
},
};
events.push(make_tagged_event("content_block_start", &ev));
}
self.tool_call_states[tc_index]
.accumulated_args
.push_str(args);
let block_index = self.tool_call_states[tc_index].block_index;
let ev = AnthropicStreamEvent::ContentBlockDelta {
index: block_index,
delta: AnthropicDelta::InputJsonDelta {
partial_json: args.clone(),
},
};
events.push(make_tagged_event("content_block_delta", &ev));
if tc.id.is_some()
&& func.name.is_some()
&& !self.tool_call_states[tc_index].stopped
{
self.tool_call_states[tc_index].stopped = true;
let ev =
AnthropicStreamEvent::ContentBlockStop { index: block_index };
events.push(make_tagged_event("content_block_stop", &ev));
}
}
}
}
}
}
events
}
fn emit_end_events_tagged(&mut self) -> Vec<TaggedEvent> {
let mut events = Vec::new();
if self.thinking_block_started && !self.thinking_block_closed {
self.thinking_block_closed = true;
let ev = AnthropicStreamEvent::ContentBlockDelta {
index: self.thinking_block_index,
delta: AnthropicDelta::SignatureDelta {
signature: "erased".to_string(),
},
};
events.push(make_tagged_event("content_block_delta", &ev));
let ev = AnthropicStreamEvent::ContentBlockStop {
index: self.thinking_block_index,
};
events.push(make_tagged_event("content_block_stop", &ev));
}
if self.text_block_started && !self.text_block_closed {
let ev = AnthropicStreamEvent::ContentBlockStop {
index: self.text_block_index,
};
events.push(make_tagged_event("content_block_stop", &ev));
}
for tc in &self.tool_call_states {
if tc.started && !tc.stopped {
let ev = AnthropicStreamEvent::ContentBlockStop {
index: tc.block_index,
};
events.push(make_tagged_event("content_block_stop", &ev));
}
}
let ev = AnthropicStreamEvent::MessageDelta {
delta: AnthropicMessageDeltaBody {
stop_reason: self.stop_reason.clone(),
stop_sequence: None,
},
usage: AnthropicUsage {
input_tokens: self.input_token_count,
output_tokens: self.output_token_count,
cache_creation_input_tokens: None,
cache_read_input_tokens: self.cached_token_count,
},
};
events.push(make_tagged_event("message_delta", &ev));
let ev = AnthropicStreamEvent::MessageStop {};
events.push(make_tagged_event("message_stop", &ev));
events
}
}
#[cfg(test)]
mod tests {
use super::*;
use dynamo_protocols::types::{
ChatChoiceStream, ChatCompletionMessageContent, ChatCompletionMessageToolCallChunk,
ChatCompletionStreamResponseDelta, FunctionCallStream, FunctionType,
};
fn text_chunk(text: &str) -> NvCreateChatCompletionStreamResponse {
#[allow(deprecated)]
NvCreateChatCompletionStreamResponse {
inner: dynamo_protocols::types::CreateChatCompletionStreamResponse {
id: "chat-1".into(),
choices: vec![ChatChoiceStream {
index: 0,
delta: ChatCompletionStreamResponseDelta {
content: Some(ChatCompletionMessageContent::Text(text.into())),
function_call: None,
tool_calls: None,
role: None,
refusal: None,
reasoning_content: None,
},
finish_reason: None,
logprobs: None,
}],
created: 0,
model: "test".into(),
service_tier: None,
system_fingerprint: None,
object: "chat.completion.chunk".into(),
usage: None,
},
nvext: None,
llm_metrics: None,
}
}
fn tool_call_chunk(
tc_index: u32,
id: Option<&str>,
name: Option<&str>,
args: Option<&str>,
) -> NvCreateChatCompletionStreamResponse {
#[allow(deprecated)]
NvCreateChatCompletionStreamResponse {
inner: dynamo_protocols::types::CreateChatCompletionStreamResponse {
id: "chat-1".into(),
choices: vec![ChatChoiceStream {
index: 0,
delta: ChatCompletionStreamResponseDelta {
content: None,
function_call: None,
tool_calls: Some(vec![ChatCompletionMessageToolCallChunk {
index: tc_index,
id: id.map(String::from),
r#type: Some(FunctionType::Function),
function: Some(FunctionCallStream {
name: name.map(String::from),
arguments: args.map(String::from),
}),
}]),
role: None,
refusal: None,
reasoning_content: None,
},
finish_reason: None,
logprobs: None,
}],
created: 0,
model: "test".into(),
service_tier: None,
system_fingerprint: None,
object: "chat.completion.chunk".into(),
usage: None,
},
nvext: None,
llm_metrics: None,
}
}
fn event_types(events: &[TaggedEvent]) -> Vec<&str> {
events.iter().map(|e| e.event_type.as_str()).collect()
}
#[test]
fn test_append_events_reuse_caller_storage() {
let mut conv = AnthropicStreamConverter::new("test-model".into(), 0);
let mut events = Vec::with_capacity(4);
conv.append_start_events(&mut events);
assert_eq!(events.len(), 1);
assert!(events.iter().all(Result::is_ok));
events.clear();
let capacity = events.capacity();
conv.append_chunk_events(&text_chunk("I'll edit the file."), &mut events);
assert_eq!(events.len(), 2);
assert_eq!(events.capacity(), capacity);
assert!(events.iter().all(Result::is_ok));
events.clear();
conv.append_chunk_events(
&tool_call_chunk(
0,
Some("call-1"),
Some("Edit"),
Some("{\"file_path\":\"/tmp/test.txt\"}"),
),
&mut events,
);
assert_eq!(events.len(), 4);
assert_eq!(events.capacity(), capacity);
assert!(events.iter().all(Result::is_ok));
events.clear();
conv.append_end_events(&mut events);
assert_eq!(events.len(), 2);
assert_eq!(events.capacity(), capacity);
assert!(events.iter().all(Result::is_ok));
events.clear();
conv.append_error_events(&mut events);
assert_eq!(events.len(), 1);
assert_eq!(events.capacity(), capacity);
assert!(events.iter().all(Result::is_ok));
}
#[test]
fn test_text_block_stops_before_tool_block_starts() {
let mut conv = AnthropicStreamConverter::new("test-model".into(), 0);
let text_events = conv.process_chunk_tagged(&text_chunk("I'll edit the file."));
assert_eq!(
event_types(&text_events),
vec!["content_block_start", "content_block_delta"]
);
let tool_events = conv.process_chunk_tagged(&tool_call_chunk(
0,
Some("call-1"),
Some("Edit"),
Some("{\"file_path\":\"/tmp/test.txt\"}"),
));
assert_eq!(
event_types(&tool_events),
vec![
"content_block_stop",
"content_block_start",
"content_block_delta",
"content_block_stop",
],
"text block must be closed before tool block starts; complete tool call stopped inline"
);
match &tool_events[0].data {
AnthropicStreamEvent::ContentBlockStop { index } => assert_eq!(*index, 0),
other => panic!("expected ContentBlockStop, got {other:?}"),
}
match &tool_events[1].data {
AnthropicStreamEvent::ContentBlockStart {
index,
content_block,
} => {
assert_eq!(*index, 1);
match content_block {
AnthropicResponseContentBlock::ToolUse { name, .. } => {
assert_eq!(name, "Edit");
}
other => panic!("expected ToolUse, got {other:?}"),
}
}
other => panic!("expected ContentBlockStart, got {other:?}"),
}
let end_events = conv.emit_end_events_tagged();
assert_eq!(
event_types(&end_events),
vec!["message_delta", "message_stop"],
"no block stops in end events (both text and tool already closed inline)"
);
}
#[test]
fn test_tool_only_response_no_text_block() {
let mut conv = AnthropicStreamConverter::new("test-model".into(), 0);
let tool_events = conv.process_chunk_tagged(&tool_call_chunk(
0,
Some("call-1"),
Some("Read"),
Some("{\"path\":\"/tmp/test.txt\"}"),
));
assert_eq!(
event_types(&tool_events),
vec![
"content_block_start",
"content_block_delta",
"content_block_stop"
],
"complete tool call emits stop inline"
);
let end_events = conv.emit_end_events_tagged();
assert_eq!(
event_types(&end_events),
vec!["message_delta", "message_stop"],
"no block stop in end events (already stopped inline)"
);
}
#[test]
fn test_text_only_response_stop_in_end_events() {
let mut conv = AnthropicStreamConverter::new("test-model".into(), 0);
conv.process_chunk_tagged(&text_chunk("Hello world"));
let end_events = conv.emit_end_events_tagged();
assert_eq!(
event_types(&end_events),
vec!["content_block_stop", "message_delta", "message_stop"]
);
match &end_events[0].data {
AnthropicStreamEvent::ContentBlockStop { index } => assert_eq!(*index, 0),
other => panic!("expected text stop at index 0, got {other:?}"),
}
}
fn reasoning_chunk(text: &str) -> NvCreateChatCompletionStreamResponse {
#[allow(deprecated)]
NvCreateChatCompletionStreamResponse {
inner: dynamo_protocols::types::CreateChatCompletionStreamResponse {
id: "chat-1".into(),
choices: vec![ChatChoiceStream {
index: 0,
delta: ChatCompletionStreamResponseDelta {
content: None,
function_call: None,
tool_calls: None,
role: None,
refusal: None,
reasoning_content: Some(text.into()),
},
finish_reason: None,
logprobs: None,
}],
created: 0,
model: "test".into(),
service_tier: None,
system_fingerprint: None,
object: "chat.completion.chunk".into(),
usage: None,
},
nvext: None,
llm_metrics: None,
}
}
#[test]
fn test_thinking_text_then_tool_call() {
let mut conv = AnthropicStreamConverter::new("test-model".into(), 0);
let ev = conv.process_chunk_tagged(&reasoning_chunk("Let me think..."));
assert_eq!(
event_types(&ev),
vec!["content_block_start", "content_block_delta"]
);
assert!(matches!(
&ev[0].data,
AnthropicStreamEvent::ContentBlockStart {
index: 0,
content_block: AnthropicResponseContentBlock::Thinking { .. }
}
));
let ev = conv.process_chunk_tagged(&text_chunk("Hello!"));
assert_eq!(
event_types(&ev),
vec![
"content_block_delta",
"content_block_stop",
"content_block_start",
"content_block_delta"
]
);
assert!(matches!(
&ev[1].data,
AnthropicStreamEvent::ContentBlockStop { index: 0 }
));
assert!(matches!(
&ev[2].data,
AnthropicStreamEvent::ContentBlockStart { index: 1, .. }
));
let ev = conv.process_chunk_tagged(&tool_call_chunk(
0,
Some("call-1"),
Some("Read"),
Some("{\"path\":\"/tmp/test.txt\"}"),
));
assert_eq!(
event_types(&ev),
vec![
"content_block_stop",
"content_block_start",
"content_block_delta",
"content_block_stop"
]
);
assert!(matches!(
&ev[0].data,
AnthropicStreamEvent::ContentBlockStop { index: 1 }
));
assert!(matches!(
&ev[1].data,
AnthropicStreamEvent::ContentBlockStart { index: 2, .. }
));
}
#[test]
fn test_thinking_only_closed_in_end_events() {
let mut conv = AnthropicStreamConverter::new("test-model".into(), 0);
conv.process_chunk_tagged(&reasoning_chunk("Deep thought..."));
let ev = conv.emit_end_events_tagged();
assert_eq!(
event_types(&ev),
vec![
"content_block_delta",
"content_block_stop",
"message_delta",
"message_stop"
]
);
}
#[test]
fn test_multiple_tool_calls_each_stopped_inline() {
let mut conv = AnthropicStreamConverter::new("test-model".into(), 0);
let events1 = conv.process_chunk_tagged(&tool_call_chunk(
0,
Some("call-1"),
Some("Read"),
Some("{\"path\":\"/tmp/a.txt\"}"),
));
assert_eq!(
event_types(&events1),
vec![
"content_block_start",
"content_block_delta",
"content_block_stop"
],
"first tool call closed inline"
);
let events2 = conv.process_chunk_tagged(&tool_call_chunk(
1,
Some("call-2"),
Some("Write"),
Some("{\"path\":\"/tmp/b.txt\"}"),
));
assert_eq!(
event_types(&events2),
vec![
"content_block_start",
"content_block_delta",
"content_block_stop"
],
"second tool call closed inline"
);
let end_events = conv.emit_end_events_tagged();
assert_eq!(
event_types(&end_events),
vec!["message_delta", "message_stop"],
"no block stops in end events"
);
}
#[test]
fn test_with_context_preserves_context() {
use crate::protocols::unified::AnthropicContext;
let ctx = AnthropicContext {
service_tier: Some("priority".to_string()),
..Default::default()
};
let mut conv = AnthropicStreamConverter::with_context("test-model".into(), 0, ctx);
assert!(conv.api_context.is_some());
assert_eq!(
conv.api_context.as_ref().unwrap().service_tier.as_deref(),
Some("priority")
);
let ev = conv.process_chunk_tagged(&text_chunk("Hello"));
assert_eq!(
event_types(&ev),
vec!["content_block_start", "content_block_delta"]
);
let end = conv.emit_end_events_tagged();
assert_eq!(
event_types(&end),
vec!["content_block_stop", "message_delta", "message_stop"]
);
}
}