1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
// Shared streaming tool-call accumulator (EVE-672)
//
// Native streaming drivers assemble tool calls from fragments spread across many
// SSE chunks: OpenAI Chat Completions keys fragments by a numeric `index`,
// OpenAI Open Responses keys them by an `item_id` string, and Gemini pushes
// whole calls. Each driver previously open-coded this into an
// `Arc<Mutex<Vec<...>>>` with subtly different growth, argument-append, and
// finalize rules. This module centralizes the accumulation so the rules live in
// one tested place and the drivers share a single `StreamToolCallAccumulator`.
//
// The accumulator keeps arguments as an un-parsed `String` fragment buffer and
// parses the JSON exactly once at finalize (EVE-636: `push_str` is amortized
// O(total), versus re-parsing a `Value` per delta which was O(n^2)). It exposes
// two finalize modes matching the historical driver behavior:
//
// - `take_finalized`: normal finish path — malformed argument JSON degrades to
// an empty object (`{}`), because the provider signalled a real tool-call
// completion.
// - `take_pending_strict`: fallback flush at `[DONE]`/end-of-stream without a
// tool-call finish — malformed argument JSON causes the call to be dropped,
// since there was no explicit completion to trust.
//
// Either fallback logs a `tracing` warning naming the call, so a degraded or
// dropped call is never silent.
use crate::tool_types::ToolCall;
use serde_json::{Value, json};
/// One in-progress tool call being assembled from streamed fragments.
#[derive(Debug, Clone)]
struct PartialToolCall {
/// Fragment key: the numeric `index` (Chat Completions) or the stream
/// `item_id` (Open Responses). Whole-call providers (Gemini) leave it empty.
key: String,
/// Provider call id, applied to the finalized [`ToolCall::id`].
id: String,
/// Function name.
name: String,
/// Accumulated JSON argument fragments (parsed once at finalize).
arguments: String,
}
/// Accumulates streamed tool-call fragments into finalized [`ToolCall`]s.
///
/// Fragments are addressed by a string `key` (the provider's per-call index or
/// item id). Calls are emitted in first-seen order.
#[derive(Debug, Default)]
pub struct StreamToolCallAccumulator {
calls: Vec<PartialToolCall>,
/// Calls discarded so far by [`Self::take_at_stream_end`] (cut-off or
/// rejected responses). Reported in completion metadata, never reset.
dropped: u32,
}
impl StreamToolCallAccumulator {
/// Create an empty accumulator.
pub fn new() -> Self {
Self::default()
}
/// Whether any fragments have been accumulated.
pub fn is_empty(&self) -> bool {
self.calls.is_empty()
}
#[expect(
clippy::expect_used,
reason = "The slot was just pushed, so the accumulator cannot be empty"
)]
fn slot(&mut self, key: &str) -> &mut PartialToolCall {
if let Some(pos) = self.calls.iter().position(|c| c.key == key) {
return &mut self.calls[pos];
}
self.calls.push(PartialToolCall {
key: key.to_string(),
id: String::new(),
name: String::new(),
arguments: String::new(),
});
self.calls.last_mut().expect("just pushed")
}
/// Apply an OpenAI Chat Completions streamed tool-call delta, keyed by the
/// chunk's numeric `index`. Any of id/name/arguments may be absent in a
/// given delta; argument fragments are appended in place.
pub fn apply_indexed_delta(
&mut self,
index: u32,
id: Option<&str>,
name: Option<&str>,
arguments: Option<&str>,
) {
let slot = self.slot(&index.to_string());
if let Some(id) = id {
slot.id = id.to_string();
}
if let Some(name) = name {
slot.name = name.to_string();
}
if let Some(args) = arguments {
slot.arguments.push_str(args);
}
}
/// Append an Open Responses `function_call_arguments.delta` fragment, keyed
/// by `item_id`. Creates the slot if the item was not yet announced.
pub fn append_arguments(&mut self, item_id: &str, delta: &str) {
self.slot(item_id).arguments.push_str(delta);
}
/// Record an Open Responses `output_item.added` function call, keyed by
/// `item_id`, setting its name and provider `call_id`.
pub fn set_item(&mut self, item_id: &str, call_id: &str, name: &str) {
let slot = self.slot(item_id);
slot.id = call_id.to_string();
slot.name = name.to_string();
}
/// Push a fully-formed tool call (Gemini emits whole `functionCall` parts
/// with already-parsed arguments, so there is nothing to reassemble).
pub fn push_complete(&mut self, id: String, name: String, arguments: Value) {
self.calls.push(PartialToolCall {
key: String::new(),
id,
name,
// Store as a compact JSON string so the single finalize path parses
// it back identically to the fragment-assembled calls.
arguments: arguments.to_string(),
});
}
/// Drain all accumulated calls, parsing each argument string once. Malformed
/// argument JSON degrades to an empty object — the normal finish path, where
/// the provider explicitly signalled tool-call completion.
pub fn take_finalized(&mut self) -> Vec<ToolCall> {
std::mem::take(&mut self.calls)
.into_iter()
.map(|c| {
let arguments = parse_arguments(&c.arguments).unwrap_or_else(|| {
warn_malformed_arguments(&c, "degraded to {}");
json!({})
});
ToolCall {
id: c.id,
name: c.name,
arguments,
}
})
.collect()
}
/// Drain all accumulated calls for a *fallback* flush (end-of-stream without
/// an explicit tool-call finish). Malformed argument JSON drops the call
/// rather than fabricating an empty object, since no completion was signalled.
pub fn take_pending_strict(&mut self) -> Vec<ToolCall> {
std::mem::take(&mut self.calls)
.into_iter()
.filter_map(|c| {
let Some(arguments) = parse_arguments(&c.arguments) else {
warn_malformed_arguments(&c, "call dropped");
return None;
};
Some(ToolCall {
id: c.id,
name: c.name,
arguments,
})
})
.collect()
}
/// Drain calls still pending when a Chat Completions stream ends.
///
/// Calls survive only when the provider omitted a finish reason or
/// reported `tool_calls`: `length` or `content_filter` mean the response
/// was cut or rejected, so pending calls are discarded rather than
/// executed. Surviving calls go through the strict flush, since no
/// explicit completion chunk vouched for their arguments.
pub fn take_at_stream_end(&mut self, finish_reason: Option<&str>) -> Vec<ToolCall> {
if !matches!(finish_reason, None | Some("tool_calls")) {
// Drain so a repeated flush cannot re-emit, but do not execute.
// Count what was discarded so the drop is observable.
self.dropped += self.discard();
return Vec::new();
}
self.take_pending_strict()
}
/// Discard every pending call without executing it, returning how many
/// there were. For responses that ended cut off or rejected.
pub fn discard(&mut self) -> u32 {
let count = u32::try_from(self.calls.len()).unwrap_or(u32::MAX);
self.calls.clear();
count
}
/// Calls [`Self::take_at_stream_end`] discarded over this stream.
pub fn dropped_at_stream_end(&self) -> u32 {
self.dropped
}
/// Drain accumulated calls, keeping only those with a non-empty name and
/// parsing arguments once (empty/malformed → `{}`). Matches the Open
/// Responses finalize path, which skips announced-but-unnamed items.
pub fn take_named(&mut self) -> Vec<ToolCall> {
std::mem::take(&mut self.calls)
.into_iter()
.filter(|c| !c.name.is_empty())
.map(|c| {
let arguments = parse_arguments(&c.arguments).unwrap_or_else(|| {
warn_malformed_arguments(&c, "degraded to {}");
json!({})
});
ToolCall {
id: c.id,
name: c.name,
arguments,
}
})
.collect()
}
}
/// Say, rather than stay silent, when streamed tool arguments did not parse.
///
/// The fallback (an empty object, or dropping the call) is the contract, but a
/// host debugging a tool loop needs to see that it happened. Only the size is
/// logged: argument text is model output and may carry user data.
fn warn_malformed_arguments(call: &PartialToolCall, outcome: &str) {
tracing::warn!(
tool_call_id = %call.id,
tool_name = %call.name,
argument_bytes = call.arguments.len(),
outcome,
"streamed tool-call arguments are not valid JSON"
);
}
/// Parse an accumulated argument fragment buffer into JSON. An empty buffer is
/// treated as an empty object; a non-empty buffer that fails to parse returns
/// `None` so callers can choose their fallback.
fn parse_arguments(buffer: &str) -> Option<Value> {
if buffer.is_empty() {
return Some(json!({}));
}
serde_json::from_str(buffer).ok()
}
#[cfg(test)]
mod tests {
use super::*;
fn calls_json(calls: Vec<ToolCall>) -> Value {
serde_json::to_value(calls).unwrap()
}
#[test]
fn interleaved_sparse_indexes_keep_first_seen_order_and_complete_payloads() {
let mut acc = StreamToolCallAccumulator::new();
acc.apply_indexed_delta(42, None, None, Some("{\"city\":"));
acc.apply_indexed_delta(2, Some("second"), Some("other"), Some("{\"n\":2}"));
acc.apply_indexed_delta(42, Some("first"), Some("weather"), Some("\"Paris 🦀\"}"));
assert_eq!(
calls_json(acc.take_finalized()),
json!([
{"id":"first","name":"weather","arguments":{"city":"Paris 🦀"}},
{"id":"second","name":"other","arguments":{"n":2}}
])
);
assert!(acc.is_empty());
assert!(acc.take_finalized().is_empty());
acc.apply_indexed_delta(42, Some("fresh"), Some("again"), Some("null"));
assert_eq!(
calls_json(acc.take_finalized()),
json!([{"id":"fresh","name":"again","arguments":null}])
);
}
#[test]
fn finalization_modes_have_distinct_malformed_and_unnamed_contracts() {
for mode in ["normal", "strict", "named"] {
let mut acc = StreamToolCallAccumulator::new();
acc.apply_indexed_delta(0, Some("empty"), Some("noop"), None);
acc.apply_indexed_delta(1, Some("bad"), Some("broken"), Some("{bad"));
acc.apply_indexed_delta(2, Some("valid"), Some("good"), Some("{\"ok\":true}"));
acc.apply_indexed_delta(3, Some("unnamed"), None, Some("[]"));
let actual = match mode {
"normal" => acc.take_finalized(),
"strict" => acc.take_pending_strict(),
_ => acc.take_named(),
};
let expected = match mode {
"normal" => json!([
{"id":"empty","name":"noop","arguments":{}}, {"id":"bad","name":"broken","arguments":{}},
{"id":"valid","name":"good","arguments":{"ok":true}}, {"id":"unnamed","name":"","arguments":[]}
]),
"strict" => json!([
{"id":"empty","name":"noop","arguments":{}},
{"id":"valid","name":"good","arguments":{"ok":true}}, {"id":"unnamed","name":"","arguments":[]}
]),
_ => json!([
{"id":"empty","name":"noop","arguments":{}}, {"id":"bad","name":"broken","arguments":{}},
{"id":"valid","name":"good","arguments":{"ok":true}}
]),
};
assert_eq!(calls_json(actual), expected, "{mode}");
assert!(acc.is_empty(), "{mode} must drain rejected entries too");
assert!(acc.take_finalized().is_empty());
}
}
#[test]
fn stream_end_counts_calls_discarded_by_a_cut_off_response() {
for (finish, dropped, survivors) in [
(Some("length"), 2, 0),
(Some("content_filter"), 2, 0),
(Some("tool_calls"), 0, 2),
(None, 0, 2),
] {
let mut acc = StreamToolCallAccumulator::new();
acc.apply_indexed_delta(0, Some("a"), Some("first"), Some("{}"));
acc.apply_indexed_delta(1, Some("b"), Some("second"), Some("{\"x\":1}"));
assert_eq!(
acc.take_at_stream_end(finish).len(),
survivors,
"{finish:?}"
);
assert_eq!(acc.dropped_at_stream_end(), dropped, "{finish:?}");
// A repeated flush neither re-emits nor double-counts.
assert!(acc.take_at_stream_end(finish).is_empty());
assert_eq!(acc.dropped_at_stream_end(), dropped, "{finish:?}");
}
}
#[test]
fn item_fragments_can_precede_metadata_without_losing_order_or_arguments() {
let mut acc = StreamToolCallAccumulator::new();
acc.append_arguments("late", "{\"city\":");
acc.set_item("second", "call_2", "second_tool");
acc.append_arguments("second", "{}");
acc.set_item("late", "call_1", "weather");
acc.append_arguments("late", "\"Paris\"}");
acc.append_arguments("orphan", "{}");
assert_eq!(
calls_json(acc.take_named()),
json!([
{"id":"call_1","name":"weather","arguments":{"city":"Paris"}},
{"id":"call_2","name":"second_tool","arguments":{}}
])
);
assert!(acc.is_empty());
}
#[test]
fn complete_calls_preserve_every_json_shape_and_distinct_empty_keys() {
let mut acc = StreamToolCallAccumulator::new();
acc.push_complete("a".into(), "first".into(), json!({"a":1,"b":"x"}));
acc.push_complete(
"b".into(),
"second".into(),
json!([null, true, 9007199254740993_u64, "🦀"]),
);
acc.push_complete("c".into(), "third".into(), Value::Null);
assert_eq!(
calls_json(acc.take_finalized()),
json!([
{"id":"a","name":"first","arguments":{"a":1,"b":"x"}},
{"id":"b","name":"second","arguments":[null,true,9007199254740993_u64,"🦀"]},
{"id":"c","name":"third","arguments":null}
])
);
assert!(acc.is_empty());
}
}