use super::buffers::StreamBuffers;
use crate::types::{
AnthropicBlockStart, AnthropicDelta, AnthropicEvent, StreamEvent, anthropic_finish_reason,
};
use crate::{Error, Result};
pub struct AnthropicAccumulator {
buffers: StreamBuffers,
}
impl AnthropicAccumulator {
pub fn new() -> Self {
Self {
buffers: StreamBuffers::new(),
}
}
pub fn capture_reasoning(mut self, capture: bool) -> Self {
self.buffers.set_capture_reasoning(capture);
self
}
pub fn process_event(&mut self, event: AnthropicEvent) -> Result<Vec<StreamEvent>> {
match event {
AnthropicEvent::ContentBlockStart {
index,
content_block,
} => {
self.open_block(index, content_block);
Ok(Vec::new())
}
AnthropicEvent::ContentBlockDelta { index, delta } => {
self.append_delta(index, delta);
Ok(Vec::new())
}
AnthropicEvent::MessageDelta { delta } => match delta.stop_reason {
Some(raw) => {
self.buffers.record_finish(anthropic_finish_reason(&raw));
self.buffers.flush()
}
None => Ok(Vec::new()),
},
AnthropicEvent::Error { error } => Err(stream_error(&error)),
AnthropicEvent::MessageStart {}
| AnthropicEvent::ContentBlockStop { .. }
| AnthropicEvent::MessageStop {}
| AnthropicEvent::Ping {}
| AnthropicEvent::Unknown => Ok(Vec::new()),
}
}
fn open_block(&mut self, index: u32, block: AnthropicBlockStart) {
match block {
AnthropicBlockStart::Text { text } => self.buffers.push_text(&text),
AnthropicBlockStart::Thinking { thinking } => self.buffers.push_reasoning(&thinking),
AnthropicBlockStart::ToolUse { id, name } => {
let call = self.buffers.tool_call(index);
call.id = Some(id);
call.name = Some(name);
}
AnthropicBlockStart::RedactedThinking {} | AnthropicBlockStart::Unknown => {}
}
}
fn append_delta(&mut self, index: u32, delta: AnthropicDelta) {
match delta {
AnthropicDelta::TextDelta { text } => self.buffers.push_text(&text),
AnthropicDelta::ThinkingDelta { thinking } => self.buffers.push_reasoning(&thinking),
AnthropicDelta::InputJsonDelta { partial_json } => {
if let Some(call) = self.buffers.open_tool_call(index) {
call.arguments.push_str(&partial_json);
}
}
AnthropicDelta::SignatureDelta {} | AnthropicDelta::Unknown => {}
}
}
pub fn finalize(&mut self) -> Result<Vec<StreamEvent>> {
self.buffers.finalize()
}
}
impl Default for AnthropicAccumulator {
fn default() -> Self {
Self::new()
}
}
fn stream_error(error: &crate::types::AnthropicErrorBody) -> Error {
let (described, status) = match error.error_type.as_deref() {
Some(kind) => {
let status = match kind {
"overloaded_error" => Some(529),
"rate_limit_error" => Some(429),
"api_error" => Some(500),
_ => None,
};
(format!("{kind}: {}", error.message), status)
}
None => (error.message.clone(), None),
};
match status {
Some(status) => Error::api_status(status, described),
None => Error::stream(described),
}
}
impl super::EventAccumulator for AnthropicAccumulator {
type Event = AnthropicEvent;
fn process(&mut self, event: Self::Event) -> Result<Vec<StreamEvent>> {
self.process_event(event)
}
fn finish(&mut self) -> Result<Vec<StreamEvent>> {
self.finalize()
}
}
#[cfg(test)]
mod tests;