use std::collections::HashMap;
#[derive(Debug, Clone, PartialEq)]
pub struct ParsedToolCall {
pub id: String,
pub name: String,
pub arguments: String,
}
#[derive(Debug, Clone, PartialEq)]
pub enum Piece {
Content(String),
Reasoning(String),
Call(ParsedToolCall),
}
enum State {
Prethink,
PostThink,
Scan,
InCall,
GemmaLabel,
GemmaThought,
}
const OPEN: &str = "<tool_call>";
const CLOSE: &str = "</tool_call>";
const THINK_END: &str = "</think>";
const GEMMA_OPEN: &str = "<|channel>";
const GEMMA_CLOSE: &str = "<channel|>";
pub struct ToolStreamParser {
state: State,
buf: String,
schemas: HashMap<String, HashMap<String, String>>,
n_calls: usize,
scan_tools: bool,
gemma: bool,
include_reasoning: bool,
postthink_nl: u8,
}
pub fn partial_suffix_len(s: &str, tag: &str) -> usize {
let max = (tag.len() - 1).min(s.len());
for k in (1..=max).rev() {
if s.ends_with(&tag[..k]) {
return k;
}
}
0
}
impl ToolStreamParser {
pub fn new(schemas: HashMap<String, HashMap<String, String>>, skip_think: bool) -> Self {
Self {
state: if skip_think { State::Prethink } else { State::Scan },
buf: String::new(),
schemas,
n_calls: 0,
scan_tools: true,
gemma: false,
include_reasoning: true,
postthink_nl: 0,
}
}
pub fn reasoning_only(include_reasoning: bool) -> Self {
let mut p = Self::new(HashMap::new(), true);
p.scan_tools = false;
p.include_reasoning = include_reasoning;
p
}
pub fn gemma_thought(include_reasoning: bool) -> Self {
let mut p = Self::new(HashMap::new(), false);
p.scan_tools = false;
p.gemma = true;
p.include_reasoning = include_reasoning;
p
}
pub fn with_include_reasoning(mut self, include: bool) -> Self {
self.include_reasoning = include;
self
}
pub fn push(&mut self, text: &str) -> Vec<Piece> {
self.buf.push_str(text);
let mut out = Vec::new();
loop {
match self.state {
State::Prethink => {
if let Some(i) = self.buf.find(THINK_END) {
self.emit_reasoning(&mut out, self.buf[..i].to_string());
self.buf.drain(..i + THINK_END.len());
self.state = State::PostThink;
self.postthink_nl = 2;
continue;
}
let keep = partial_suffix_len(&self.buf, THINK_END);
let emit_to = self.buf.len() - keep;
if emit_to > 0 {
self.emit_reasoning(&mut out, self.buf[..emit_to].to_string());
self.buf.drain(..emit_to);
}
break;
}
State::PostThink => {
while self.postthink_nl > 0 && self.buf.starts_with('\n') {
self.buf.drain(..1);
self.postthink_nl -= 1;
}
if self.postthink_nl > 0 && self.buf.is_empty() {
break; }
self.state = State::Scan;
continue;
}
State::Scan => {
if self.gemma {
if let Some(i) = self.buf.find(GEMMA_OPEN) {
if i > 0 {
emit_content(&mut out, self.buf[..i].to_string());
}
self.buf.drain(..i + GEMMA_OPEN.len());
self.state = State::GemmaLabel;
continue;
}
let keep = partial_suffix_len(&self.buf, GEMMA_OPEN);
let emit_to = self.buf.len() - keep;
if emit_to > 0 {
emit_content(&mut out, self.buf[..emit_to].to_string());
self.buf.drain(..emit_to);
}
break;
}
if !self.scan_tools {
if !self.buf.is_empty() {
emit_content(&mut out, std::mem::take(&mut self.buf));
}
break;
}
if let Some(i) = self.buf.find(OPEN) {
if i > 0 {
emit_content(&mut out, self.buf[..i].to_string());
}
self.buf.drain(..i + OPEN.len());
self.state = State::InCall;
continue;
}
let keep = partial_suffix_len(&self.buf, OPEN);
let emit_to = self.buf.len() - keep;
if emit_to > 0 {
emit_content(&mut out, self.buf[..emit_to].to_string());
self.buf.drain(..emit_to);
}
break;
}
State::InCall => {
let Some(i) = self.buf.find(CLOSE) else { break };
let inner: String = self.buf[..i].to_string();
self.buf.drain(..i + CLOSE.len());
self.state = State::Scan;
match self.parse_block(&inner) {
Some(call) => out.push(Piece::Call(call)),
None => emit_content(&mut out, format!("{OPEN}{inner}{CLOSE}")),
}
continue;
}
State::GemmaLabel => {
let Some(i) = self.buf.find('\n') else { break };
self.buf.drain(..i + 1);
self.state = State::GemmaThought;
continue;
}
State::GemmaThought => {
if let Some(i) = self.buf.find(GEMMA_CLOSE) {
let text = self.buf[..i].strip_suffix('\n').unwrap_or(&self.buf[..i]);
self.emit_reasoning(&mut out, text.to_string());
self.buf.drain(..i + GEMMA_CLOSE.len());
self.state = State::Scan;
continue;
}
let mut keep = partial_suffix_len(&self.buf, GEMMA_CLOSE);
if self.buf[..self.buf.len() - keep].ends_with('\n') {
keep += 1;
}
let emit_to = self.buf.len() - keep;
if emit_to > 0 {
self.emit_reasoning(&mut out, self.buf[..emit_to].to_string());
self.buf.drain(..emit_to);
}
break;
}
}
}
out
}
pub fn finish(&mut self) -> Vec<Piece> {
let mut out = Vec::new();
if !self.buf.is_empty() {
let tail = std::mem::take(&mut self.buf);
match self.state {
State::Prethink => self.emit_reasoning(&mut out, tail),
State::InCall => emit_content(&mut out, format!("{OPEN}{tail}")),
State::GemmaThought => {
let t = tail.strip_suffix('\n').unwrap_or(&tail);
self.emit_reasoning(&mut out, t.to_string());
}
State::GemmaLabel => {}
_ => emit_content(&mut out, tail),
}
}
self.state = State::Scan;
out
}
pub fn n_calls(&self) -> usize {
self.n_calls
}
fn emit_reasoning(&self, out: &mut Vec<Piece>, text: String) {
if !self.include_reasoning || text.is_empty() {
return;
}
if let Some(Piece::Reasoning(prev)) = out.last_mut() {
prev.push_str(&text);
return;
}
out.push(Piece::Reasoning(text));
}
fn parse_block(&mut self, inner: &str) -> Option<ParsedToolCall> {
let s = inner.trim();
let rest = s.strip_prefix("<function=")?;
let gt = rest.find('>')?;
let name = &rest[..gt];
if name.is_empty() || name.contains(['<', '>', '\n']) {
return None;
}
let mut body = rest[gt + 1..].strip_suffix("</function>")?;
let mut args = serde_json::Map::new();
loop {
let t = body.trim_start();
if t.is_empty() {
break;
}
let r = t.strip_prefix("<parameter=")?;
let gt = r.find('>')?;
let key = &r[..gt];
if key.is_empty() || key.contains(['<', '>', '\n']) {
return None;
}
let after = &r[gt + 1..];
let after = after.strip_prefix('\n').unwrap_or(after);
let end = after.find("</parameter>")?;
let raw = after[..end].strip_suffix('\n').unwrap_or(&after[..end]);
args.insert(key.to_string(), self.coerce(name, key, raw));
body = &after[end + "</parameter>".len()..];
}
let arguments = serde_json::to_string(&serde_json::Value::Object(args)).ok()?;
let id = format!("call_{:016x}", fnv1a64(&[
&self.n_calls.to_le_bytes(), name.as_bytes(), arguments.as_bytes(),
]));
self.n_calls += 1;
Some(ParsedToolCall { id, name: name.to_string(), arguments })
}
fn coerce(&self, func: &str, param: &str, raw: &str) -> serde_json::Value {
let declared = self.schemas.get(func).and_then(|m| m.get(param)).map(String::as_str);
match declared {
Some("string") | None => serde_json::Value::String(raw.to_string()),
Some(_) => serde_json::from_str::<serde_json::Value>(raw.trim())
.unwrap_or_else(|_| serde_json::Value::String(raw.to_string())),
}
}
}
fn emit_content(out: &mut Vec<Piece>, text: String) {
if text.is_empty() {
return;
}
if let Some(Piece::Content(prev)) = out.last_mut() {
prev.push_str(&text);
return;
}
out.push(Piece::Content(text));
}
fn fnv1a64(parts: &[&[u8]]) -> u64 {
let mut h: u64 = 0xcbf29ce484222325;
for part in parts {
for &b in *part {
h ^= b as u64;
h = h.wrapping_mul(0x100000001b3);
}
}
h
}
#[cfg(test)]
mod tests {
use super::*;
fn weather_schema() -> HashMap<String, HashMap<String, String>> {
let mut params = HashMap::new();
params.insert("city".to_string(), "string".to_string());
params.insert("days".to_string(), "integer".to_string());
params.insert("metric".to_string(), "boolean".to_string());
let mut m = HashMap::new();
m.insert("get_weather".to_string(), params);
m
}
const EMISSION: &str = "I'll check.\n\n<tool_call>\n<function=get_weather>\n<parameter=city>\n\
Paris\n</parameter>\n<parameter=days>\n3\n</parameter>\n<parameter=metric>\ntrue\n</parameter>\n\
</function>\n</tool_call>";
fn reassemble(pieces: &[Piece]) -> (String, Vec<ParsedToolCall>) {
let (content, reasoning, calls) = reassemble3(pieces);
assert!(reasoning.is_empty(), "unexpected reasoning: {reasoning:?}");
(content, calls)
}
fn reassemble3(pieces: &[Piece]) -> (String, String, Vec<ParsedToolCall>) {
let mut content = String::new();
let mut reasoning = String::new();
let mut calls = Vec::new();
for p in pieces {
match p {
Piece::Content(t) => content.push_str(t),
Piece::Reasoning(t) => reasoning.push_str(t),
Piece::Call(c) => calls.push(c.clone()),
}
}
(content, reasoning, calls)
}
#[test]
fn parses_call_with_schema_coercion() {
let mut p = ToolStreamParser::new(weather_schema(), false);
let mut pieces = p.push(EMISSION);
pieces.extend(p.finish());
let (content, calls) = reassemble(&pieces);
assert_eq!(content, "I'll check.\n\n");
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].name, "get_weather");
assert_eq!(calls[0].arguments, r#"{"city":"Paris","days":3,"metric":true}"#);
assert!(calls[0].id.starts_with("call_"));
}
#[test]
fn char_by_char_deltas_produce_the_same_result() {
let mut p = ToolStreamParser::new(weather_schema(), false);
let mut pieces: Vec<Piece> = Vec::new();
for ch in EMISSION.chars() {
pieces.extend(p.push(&ch.to_string()));
}
pieces.extend(p.finish());
let (content, calls) = reassemble(&pieces);
assert_eq!(content, "I'll check.\n\n");
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].arguments, r#"{"city":"Paris","days":3,"metric":true}"#);
}
#[test]
fn think_gate_routes_think_text_to_reasoning_not_content() {
let mut p = ToolStreamParser::new(weather_schema(), true);
let text = "planning a <tool_call> here...</think>\n\n<tool_call>\n\
<function=get_weather>\n<parameter=city>\nOslo\n</parameter>\n</function>\n</tool_call>";
let mut pieces = p.push(text);
pieces.extend(p.finish());
let (content, reasoning, calls) = reassemble3(&pieces);
assert_eq!(reasoning, "planning a <tool_call> here...");
assert_eq!(content, "");
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].arguments, r#"{"city":"Oslo"}"#);
}
#[test]
fn reasoning_only_mode_splits_think_from_content_char_by_char() {
let text = "step one\nstep two</think>\n\nAnswer with a <tool_call> literal.";
for chunked in [false, true] {
let mut p = ToolStreamParser::reasoning_only(true);
let mut pieces = Vec::new();
if chunked {
for ch in text.chars() {
pieces.extend(p.push(&ch.to_string()));
}
} else {
pieces.extend(p.push(text));
}
pieces.extend(p.finish());
let (content, reasoning, calls) = reassemble3(&pieces);
assert_eq!(reasoning, "step one\nstep two", "chunked={chunked}");
assert_eq!(content, "Answer with a <tool_call> literal.", "chunked={chunked}");
assert!(calls.is_empty());
}
}
#[test]
fn include_reasoning_false_drops_think_text() {
let mut p = ToolStreamParser::reasoning_only(false);
let mut pieces = p.push("hidden plan</think>\n\nvisible answer");
pieces.extend(p.finish());
let (content, reasoning, calls) = reassemble3(&pieces);
assert_eq!(reasoning, "");
assert_eq!(content, "visible answer");
assert!(calls.is_empty());
}
#[test]
fn unclosed_think_flushes_as_reasoning() {
let mut p = ToolStreamParser::reasoning_only(true);
let mut pieces = p.push("half a thought");
pieces.extend(p.finish());
let (content, reasoning, _) = reassemble3(&pieces);
assert_eq!(reasoning, "half a thought");
assert_eq!(content, "");
}
#[test]
fn malformed_block_is_surfaced_verbatim() {
let text = "<tool_call>\n{\"name\": \"get_weather\", \"arguments\": {broken\n</tool_call>done";
let mut p = ToolStreamParser::new(weather_schema(), false);
let mut pieces = p.push(text);
pieces.extend(p.finish());
let (content, calls) = reassemble(&pieces);
assert_eq!(content, text); assert!(calls.is_empty());
}
#[test]
fn unterminated_block_flushes_raw_on_finish() {
let mut p = ToolStreamParser::new(weather_schema(), false);
let mut pieces = p.push("<tool_call>\n<function=get_weather>\n<parameter=city>\nParis");
pieces.extend(p.finish());
let (content, calls) = reassemble(&pieces);
assert_eq!(content, "<tool_call>\n<function=get_weather>\n<parameter=city>\nParis");
assert!(calls.is_empty());
}
#[test]
fn two_calls_and_multiline_string_values() {
let text = "<tool_call>\n<function=get_weather>\n<parameter=city>\nline one\nline two\n\
</parameter>\n</function>\n</tool_call>\n<tool_call>\n<function=get_weather>\n<parameter=days>\n\
not-a-number\n</parameter>\n</function>\n</tool_call>";
let mut p = ToolStreamParser::new(weather_schema(), false);
let mut pieces = p.push(text);
pieces.extend(p.finish());
let (content, calls) = reassemble(&pieces);
assert_eq!(content, "\n"); assert_eq!(calls.len(), 2);
assert_eq!(calls[0].arguments, r#"{"city":"line one\nline two"}"#);
assert_eq!(calls[1].arguments, r#"{"days":"not-a-number"}"#);
assert_ne!(calls[0].id, calls[1].id);
}
#[test]
fn gemma_thought_channel_splits_reasoning_from_content_char_by_char() {
let text = "<|channel>thought\nThe user wants ok.\nSo reply ok.\n<channel|>ok";
for chunked in [false, true] {
let mut p = ToolStreamParser::gemma_thought(true);
let mut pieces = Vec::new();
if chunked {
for ch in text.chars() {
pieces.extend(p.push(&ch.to_string()));
}
} else {
pieces.extend(p.push(text));
}
pieces.extend(p.finish());
let (content, reasoning, calls) = reassemble3(&pieces);
assert_eq!(reasoning, "The user wants ok.\nSo reply ok.", "chunked={chunked}");
assert_eq!(content, "ok", "chunked={chunked}");
assert!(calls.is_empty());
}
}
#[test]
fn gemma_content_before_and_between_channels() {
let text = "ok<|channel>thought\nreconsidering\n<channel|> more";
let mut p = ToolStreamParser::gemma_thought(true);
let mut pieces = p.push(text);
pieces.extend(p.finish());
let (content, reasoning, _) = reassemble3(&pieces);
assert_eq!(reasoning, "reconsidering");
assert_eq!(content, "ok more");
}
#[test]
fn gemma_unclosed_thought_flushes_as_reasoning_and_excludes_reasoning_drops() {
let mut p = ToolStreamParser::gemma_thought(true);
let mut pieces = p.push("<|channel>thought\nhalf a tho");
pieces.extend(p.finish());
let (content, reasoning, _) = reassemble3(&pieces);
assert_eq!(reasoning, "half a tho");
assert_eq!(content, "");
let mut p = ToolStreamParser::gemma_thought(false);
let mut pieces = p.push("<|channel>thought\nhidden\n<channel|>visible");
pieces.extend(p.finish());
let (content, reasoning, _) = reassemble3(&pieces);
assert_eq!(reasoning, "");
assert_eq!(content, "visible");
}
#[test]
fn gemma_partial_open_tag_holdback_never_loses_bytes() {
let mut p = ToolStreamParser::gemma_thought(true);
let mut pieces = p.push("a <|chan");
pieces.extend(p.push("nel of prose"));
pieces.extend(p.finish());
let (content, reasoning, _) = reassemble3(&pieces);
assert_eq!(content, "a <|channel of prose");
assert_eq!(reasoning, "");
}
#[test]
fn partial_tag_holdback_never_loses_bytes() {
let mut p = ToolStreamParser::new(HashMap::new(), false);
let mut pieces = p.push("a <tool");
pieces.extend(p.push("box holds bytes"));
pieces.extend(p.finish());
let (content, calls) = reassemble(&pieces);
assert_eq!(content, "a <toolbox holds bytes");
assert!(calls.is_empty());
}
}