feat(streaming): stream assistant tokens on-device and over the cloud
Thread an optional TokenSink through the InferenceBackend trait so a turn streams text as it is generated while still returning the full result. - LocalBackend streams via onde stream_message / stream_tool_results when no tools are offered. onde's tool-aware path must buffer to detect tool calls, so streaming covers the tools-disabled and final forced-text rounds. - OpenAiBackend streams via SSE (stream:true), reassembling content and index-keyed tool_calls; adds the reqwest "stream" feature. - TUI renders live tokens with <think> reasoning hidden. - ACP emits AgentMessageChunk deltas live, stripping <think> across chunk boundaries. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
paydii committed
Jun 27, 2026 at 08:38 UTC
33432e5e59517a4bf5309af30839d67fae8fd3e6
4 files changed
+527
-71
Cargo.toml
+1
-1
@@ -43,6 +43,6 @@ serde_json = "1"
43
toml = "0.8"
44
async-trait = "0.1"
45
regex = "1"
46
-reqwest = { version = "0.12", default-features = false, features = ["blocking", "json", "rustls-tls"] }
46
+reqwest = { version = "0.12", default-features = false, features = ["blocking", "json", "rustls-tls", "stream"] }
47
uuid = { version = "1", features = ["v4"] }
48
rpassword = "7"
src/backend.rs
+311
-6
@@ -62,24 +62,40 @@ pub struct TurnResult {
62
/// Backend errors are plain strings. Callers map them to ACP errors.
63
pub type BackendError = String;
64
65
+/// A sink for streaming assistant text deltas to the UI as they are produced.
66
+///
67
+/// When a caller passes `Some(sink)`, a streaming-capable backend forwards each
68
+/// text fragment through it as the model emits it; the returned [`TurnResult`]
69
+/// still carries the fully assembled text (and any tool calls). When the sink is
70
+/// `None`, the backend runs in non-streaming mode. Unbounded so the inference
71
+/// task never blocks on a slow consumer.
72
+pub type TokenSink = tokio::sync::mpsc::UnboundedSender<String>;
73
+
74
// ── The trait ───────────────────────────────────────────────────────────────────
75
76
/// A swappable inference backend driving siGit Code's agent loop.
77
#[async_trait]
78
pub trait InferenceBackend: Send + Sync {
79
/// Start an assistant turn from a new user message, offering `tools`.
80
+ ///
81
+ /// If `sink` is `Some`, text is streamed through it as it is generated. A
82
+ /// backend may decline to stream a given round (for example, on-device
83
+ /// inference cannot stream while it is still deciding whether to call a
84
+ /// tool); in that case the text is delivered only via the returned result.
85
async fn send_message_with_tools(
86
&self,
87
text: &str,
88
tools: &[ToolSpec],
89
+ sink: Option<&TokenSink>,
90
) -> Result<TurnResult, BackendError>;
91
92
/// Continue the turn by returning tool results. `tools` may be `None` on the
78
- /// final round to force a text answer.
93
+ /// final round to force a text answer. `sink` streams that text when set.
94
async fn send_tool_results(
95
&self,
96
results: Vec<ToolResult>,
97
tools: Option<&[ToolSpec]>,
98
+ sink: Option<&TokenSink>,
99
) -> Result<TurnResult, BackendError>;
100
101
/// Whether inference runs over the network (a configured provider) rather
@@ -118,7 +134,22 @@ impl InferenceBackend for LocalBackend {
134
&self,
135
text: &str,
136
tools: &[ToolSpec],
137
+ sink: Option<&TokenSink>,
138
) -> Result<TurnResult, BackendError> {
139
+ // onde's tool-aware path is non-streaming: it has to buffer the whole
140
+ // reply to detect tool calls. We can only stream when no tools are on
141
+ // offer (a plain answer), which is exactly the tools-disabled case.
142
+ if let Some(sink) = sink
143
+ && tools.is_empty()
144
+ {
145
+ let rx = self
146
+ .engine
147
+ .stream_message(text)
148
+ .await
149
+ .map_err(|error| error.to_string())?;
150
+ return drain_onde_stream(rx, sink).await;
151
+ }
152
+
153
let onde_tools = to_onde_tools(tools);
154
let result = self
155
.engine
@@ -132,6 +163,7 @@ impl InferenceBackend for LocalBackend {
163
&self,
164
results: Vec<ToolResult>,
165
tools: Option<&[ToolSpec]>,
166
+ sink: Option<&TokenSink>,
167
) -> Result<TurnResult, BackendError> {
168
let onde_results: Vec<onde::inference::ToolResult> = results
169
.into_iter()
@@ -140,6 +172,20 @@ impl InferenceBackend for LocalBackend {
172
content: result.content,
173
})
174
.collect();
175
+
176
+ // The final round passes `tools = None` to force a text answer; that's
177
+ // the only round onde can stream, since no further tool calls are parsed.
178
+ if let Some(sink) = sink
179
+ && tools.is_none()
180
+ {
181
+ let rx = self
182
+ .engine
183
+ .stream_tool_results(onde_results, None)
184
+ .await
185
+ .map_err(|error| error.to_string())?;
186
+ return drain_onde_stream(rx, sink).await;
187
+ }
188
+
189
let onde_tools = tools.map(to_onde_tools);
190
let result = self
191
.engine
@@ -154,6 +200,38 @@ impl InferenceBackend for LocalBackend {
200
}
201
}
202
203
+/// Drain an onde streaming receiver, forwarding each token to `sink` and
204
+/// assembling the full text. onde reports stream failures as a final chunk whose
205
+/// `finish_reason` is `"error: …"`; surface those as a backend error.
206
+async fn drain_onde_stream(
207
+ mut rx: tokio::sync::mpsc::Receiver<onde::inference::StreamChunk>,
208
+ sink: &TokenSink,
209
+) -> Result<TurnResult, BackendError> {
210
+ let mut text = String::new();
211
+ while let Some(chunk) = rx.recv().await {
212
+ if !chunk.delta.is_empty() {
213
+ text.push_str(&chunk.delta);
214
+ // The receiver is the UI; if it's gone the turn is being cancelled,
215
+ // so stop assembling rather than spinning the model to completion.
216
+ if sink.send(chunk.delta).is_err() {
217
+ break;
218
+ }
219
+ }
220
+ if chunk.done {
221
+ if let Some(reason) = chunk.finish_reason
222
+ && let Some(message) = reason.strip_prefix("error: ")
223
+ {
224
+ return Err(message.to_string());
225
+ }
226
+ break;
227
+ }
228
+ }
229
+ Ok(TurnResult {
230
+ text,
231
+ tool_calls: Vec::new(),
232
+ })
233
+}
234
+
235
/// Convert an `onde` tool-aware result into the neutral [`TurnResult`].
236
fn onde_result_to_turn(result: onde::inference::ToolAwareResult) -> TurnResult {
237
TurnResult {
@@ -230,14 +308,20 @@ impl OpenAiBackend {
308
}
309
310
/// POST the current history (plus `tools`) and apply the assistant reply to
233
- /// history, returning the neutral turn result.
234
- async fn complete(&self, tools: Option<&[ToolSpec]>) -> Result<TurnResult, BackendError> {
311
+ /// history, returning the neutral turn result. Streams via SSE when `sink`
312
+ /// is set; otherwise reads a single JSON response.
313
+ async fn complete(
314
+ &self,
315
+ tools: Option<&[ToolSpec]>,
316
+ sink: Option<&TokenSink>,
317
+ ) -> Result<TurnResult, BackendError> {
318
let url = format!("{}/chat/completions", self.base_url.trim_end_matches('/'));
319
+ let streaming = sink.is_some();
320
321
let mut body = serde_json::json!({
322
"model": self.model,
323
"messages": *self.history.lock().await,
240
- "stream": false,
324
+ "stream": streaming,
325
});
326
if let Some(tools) = tools
327
&& !tools.is_empty()
@@ -260,6 +344,15 @@ impl OpenAiBackend {
344
return Err(format!("endpoint returned {status}: {detail}"));
345
}
346
347
+ if let Some(sink) = sink {
348
+ self.consume_stream(response, sink).await
349
+ } else {
350
+ self.consume_json(response).await
351
+ }
352
+ }
353
+
354
+ /// Parse a single non-streaming chat-completion response.
355
+ async fn consume_json(&self, response: reqwest::Response) -> Result<TurnResult, BackendError> {
356
let parsed: ChatCompletion = response
357
.json()
358
.await
@@ -289,6 +382,150 @@ impl OpenAiBackend {
382
383
Ok(TurnResult { text, tool_calls })
384
}
385
+
386
+ /// Consume an OpenAI Server-Sent Events stream, forwarding content deltas to
387
+ /// `sink` and reassembling any tool calls (which arrive fragmented across
388
+ /// chunks, keyed by `index`).
389
+ async fn consume_stream(
390
+ &self,
391
+ response: reqwest::Response,
392
+ sink: &TokenSink,
393
+ ) -> Result<TurnResult, BackendError> {
394
+ use futures::StreamExt;
395
+
396
+ let mut stream = response.bytes_stream();
397
+ // Newlines are ASCII, so splitting raw bytes on `\n` never bisects a
398
+ // multibyte UTF-8 sequence; we only lossily decode whole lines.
399
+ let mut buffer: Vec<u8> = Vec::new();
400
+ let mut text = String::new();
401
+ let mut tool_accum: Vec<StreamingToolCall> = Vec::new();
402
+ let mut done = false;
403
+
404
+ while let Some(item) = stream.next().await {
405
+ let bytes = item.map_err(|error| format!("stream read error: {error}"))?;
406
+ buffer.extend_from_slice(&bytes);
407
+
408
+ while let Some(pos) = buffer.iter().position(|&b| b == b'\n') {
409
+ let line: Vec<u8> = buffer.drain(..=pos).collect();
410
+ let line = String::from_utf8_lossy(&line);
411
+ let line = line.trim();
412
+
413
+ let Some(data) = line.strip_prefix("data:") else {
414
+ continue;
415
+ };
416
+ let data = data.trim();
417
+ if data == "[DONE]" {
418
+ done = true;
419
+ break;
420
+ }
421
+ if data.is_empty() {
422
+ continue;
423
+ }
424
+
425
+ let chunk: StreamCompletion = match serde_json::from_str(data) {
426
+ Ok(chunk) => chunk,
427
+ // Skip keep-alive comments and anything we can't parse rather
428
+ // than aborting a turn over one malformed frame.
429
+ Err(_) => continue,
430
+ };
431
+
432
+ let Some(choice) = chunk.choices.into_iter().next() else {
433
+ continue;
434
+ };
435
+ if let Some(content) = choice.delta.content
436
+ && !content.is_empty()
437
+ {
438
+ text.push_str(&content);
439
+ if sink.send(content).is_err() {
440
+ // Consumer dropped (turn cancelled) — stop reading.
441
+ done = true;
442
+ break;
443
+ }
444
+ }
445
+ for delta in choice.delta.tool_calls.into_iter().flatten() {
446
+ let index = delta.index.unwrap_or(0) as usize;
447
+ if tool_accum.len() <= index {
448
+ tool_accum.resize_with(index + 1, StreamingToolCall::default);
449
+ }
450
+ let slot = &mut tool_accum[index];
451
+ if let Some(id) = delta.id {
452
+ slot.id = id;
453
+ }
454
+ if let Some(function) = delta.function {
455
+ if let Some(name) = function.name {
456
+ slot.name = name;
457
+ }
458
+ if let Some(arguments) = function.arguments {
459
+ slot.arguments.push_str(&arguments);
460
+ }
461
+ }
462
+ }
463
+ }
464
+
465
+ if done {
466
+ break;
467
+ }
468
+ }
469
+
470
+ let tool_calls: Vec<ToolCall> = tool_accum
471
+ .iter()
472
+ .filter(|call| !call.name.is_empty())
473
+ .enumerate()
474
+ .map(|(index, call)| ToolCall {
475
+ id: if call.id.is_empty() {
476
+ format!("call_{index}")
477
+ } else {
478
+ call.id.clone()
479
+ },
480
+ name: call.name.clone(),
481
+ arguments: call.arguments.clone(),
482
+ })
483
+ .collect();
484
+
485
+ // Record the assistant turn so later tool results have context.
486
+ self.history
487
+ .lock()
488
+ .await
489
+ .push(streamed_assistant_history(&text, &tool_calls));
490
+
491
+ Ok(TurnResult { text, tool_calls })
492
+ }
493
+}
494
+
495
+/// One tool call being reassembled from streamed deltas.
496
+#[derive(Default)]
497
+struct StreamingToolCall {
498
+ id: String,
499
+ name: String,
500
+ arguments: String,
501
+}
502
+
503
+/// Rebuild the assistant message for replay in history after a streamed turn,
504
+/// preserving any tool calls so the follow-up request is well-formed. Mirrors
505
+/// [`ResponseMessage::into_history_value`] for the non-streaming path.
506
+fn streamed_assistant_history(text: &str, tool_calls: &[ToolCall]) -> serde_json::Value {
507
+ let mut message = serde_json::json!({ "role": "assistant" });
508
+ message["content"] = if text.is_empty() {
509
+ serde_json::Value::Null
510
+ } else {
511
+ serde_json::Value::String(text.to_string())
512
+ };
513
+ if !tool_calls.is_empty() {
514
+ message["tool_calls"] = serde_json::json!(
515
+ tool_calls
516
+ .iter()
517
+ .map(|call| serde_json::json!({
518
+ "id": call.id,
519
+ "type": "function",
520
+ "function": {
521
+ "name": call.name,
522
+ "arguments": call.arguments,
523
+ }
524
+ }))
525
+ .collect::<Vec<_>>()
526
+ );
527
+ }
528
+ message
529
}
530
531
#[async_trait]
@@ -297,18 +534,20 @@ impl InferenceBackend for OpenAiBackend {
534
&self,
535
text: &str,
536
tools: &[ToolSpec],
537
+ sink: Option<&TokenSink>,
538
) -> Result<TurnResult, BackendError> {
539
self.history
540
.lock()
541
.await
542
.push(serde_json::json!({ "role": "user", "content": text }));
305
- self.complete(Some(tools)).await
543
+ self.complete(Some(tools), sink).await
544
}
545
546
async fn send_tool_results(
547
&self,
548
results: Vec<ToolResult>,
549
tools: Option<&[ToolSpec]>,
550
+ sink: Option<&TokenSink>,
551
) -> Result<TurnResult, BackendError> {
552
{
553
let mut history = self.history.lock().await;
@@ -320,7 +559,7 @@ impl InferenceBackend for OpenAiBackend {
559
}));
560
}
561
}
323
- self.complete(tools).await
562
+ self.complete(tools, sink).await
563
}
564
565
fn is_remote(&self) -> bool {
@@ -390,6 +629,46 @@ struct ResponseFunction {
629
arguments: String,
630
}
631
632
+// ── OpenAI streaming (SSE) chunk shapes ─────────────────────────────────────────
633
+
634
+#[derive(Debug, Deserialize)]
635
+struct StreamCompletion {
636
+ #[serde(default)]
637
+ choices: Vec<StreamChoice>,
638
+}
639
+
640
+#[derive(Debug, Deserialize)]
641
+struct StreamChoice {
642
+ #[serde(default)]
643
+ delta: StreamDelta,
644
+}
645
+
646
+#[derive(Debug, Default, Deserialize)]
647
+struct StreamDelta {
648
+ #[serde(default)]
649
+ content: Option<String>,
650
+ #[serde(default)]
651
+ tool_calls: Option<Vec<StreamToolCallDelta>>,
652
+}
653
+
654
+#[derive(Debug, Deserialize)]
655
+struct StreamToolCallDelta {
656
+ #[serde(default)]
657
+ index: Option<u32>,
658
+ #[serde(default)]
659
+ id: Option<String>,
660
+ #[serde(default)]
661
+ function: Option<StreamFunctionDelta>,
662
+}
663
+
664
+#[derive(Debug, Deserialize)]
665
+struct StreamFunctionDelta {
666
+ #[serde(default)]
667
+ name: Option<String>,
668
+ #[serde(default)]
669
+ arguments: Option<String>,
670
+}
671
+
672
#[cfg(test)]
673
mod tests {
674
use super::*;
@@ -422,6 +701,32 @@ mod tests {
701
assert_eq!(json[0]["function"]["parameters"]["type"], "object");
702
}
703
704
+ #[test]
705
+ fn streamed_assistant_history_omits_empty_tool_calls() {
706
+ let value = streamed_assistant_history("hello", &[]);
707
+ assert_eq!(value["role"], "assistant");
708
+ assert_eq!(value["content"], "hello");
709
+ assert!(value.get("tool_calls").is_none());
710
+ }
711
+
712
+ #[test]
713
+ fn streamed_assistant_history_preserves_tool_calls() {
714
+ let calls = vec![ToolCall {
715
+ id: "call_0".to_string(),
716
+ name: "read_file".to_string(),
717
+ arguments: r#"{"path":"a.rs"}"#.to_string(),
718
+ }];
719
+ let value = streamed_assistant_history("", &calls);
720
+ assert!(value["content"].is_null());
721
+ assert_eq!(value["tool_calls"][0]["id"], "call_0");
722
+ assert_eq!(value["tool_calls"][0]["type"], "function");
723
+ assert_eq!(value["tool_calls"][0]["function"]["name"], "read_file");
724
+ assert_eq!(
725
+ value["tool_calls"][0]["function"]["arguments"],
726
+ r#"{"path":"a.rs"}"#
727
+ );
728
+ }
729
+
730
#[test]
731
fn assistant_message_with_tool_calls_round_trips() {
732
let message = ResponseMessage {
src/chat.rs
+89
-35
@@ -75,7 +75,7 @@ mod tui {
75
use anyhow::Result;
76
use crossterm::event::{Event, EventStream, KeyCode, KeyEvent, KeyEventKind, KeyModifiers};
77
use futures::StreamExt;
78
- use onde::inference::{ChatEngine, SamplingConfig, StreamChunk};
78
+ use onde::inference::{ChatEngine, SamplingConfig};
79
80
use crate::backend::{InferenceBackend, LocalBackend, OpenAiBackend, ToolResult, ToolSpec};
81
use crate::models::{ModelCacheHealth, ModelPickerItem, ModelSource, build_model_picker_items};
@@ -148,6 +148,11 @@ mod tui {
148
enum InferenceUpdate {
149
/// show tool name in chat while it runs
150
ToolUse(String),
151
+ /// a streamed token fragment of the assistant's reply
152
+ Delta(String),
153
+ /// the streamed reply is complete; commit the accumulated buffer
154
+ StreamEnd,
155
+ /// a complete (non-streamed) assistant reply
156
Response(String),
157
Error(String),
158
}
@@ -164,7 +169,8 @@ mod tui {
169
input: String,
170
cursor: usize,
171
scroll_offset: u16,
167
- stream_rx: Option<mpsc::Receiver<StreamChunk>>,
172
+ /// true while assistant tokens are streaming into `stream_buf`
173
+ streaming: bool,
174
stream_buf: String,
175
inference_rx: Option<mpsc::Receiver<InferenceUpdate>>,
176
model_load_rx: Option<mpsc::Receiver<ModelLoadUpdate>>,
@@ -260,7 +266,7 @@ mod tui {
266
input: String::new(),
267
cursor: 0,
268
scroll_offset: 0,
263
- stream_rx: None,
269
+ streaming: false,
270
stream_buf: String::new(),
271
inference_rx: None,
272
model_load_rx: None,
@@ -298,11 +304,11 @@ mod tui {
304
}
305
306
fn is_streaming(&self) -> bool {
301
- self.stream_rx.is_some()
307
+ self.streaming
308
}
309
310
fn finalize_stream(&mut self) {
305
- self.stream_rx = None;
311
+ self.streaming = false;
312
if !self.stream_buf.is_empty() {
313
let text = std::mem::take(&mut self.stream_buf);
314
self.messages.push(ChatMessage::assistant(text));
@@ -311,11 +317,23 @@ mod tui {
317
}
318
319
fn push_stream_delta(&mut self, delta: &str) {
320
+ self.streaming = true;
321
self.stream_buf.push_str(delta);
322
+ // Hide reasoning the way the rest of the app does: keep the "thinking"
323
+ // spinner until visible (non-<think>) text appears, then show the
324
+ // live reply. Don't call stop_thinking() — that drops the channel.
325
+ let (_think, visible) = super::strip_think_blocks(&self.stream_buf);
326
+ self.thinking = visible.trim().is_empty();
327
self.blink_counter = self.blink_counter.wrapping_add(1);
328
self.blink_on = self.blink_counter % 4 < 2;
329
}
330
331
+ /// The portion of the streaming buffer to show live, with reasoning hidden.
332
+ fn visible_stream(&self) -> String {
333
+ let (_think, visible) = super::strip_think_blocks(&self.stream_buf);
334
+ visible
335
+ }
336
+
337
fn start_thinking(&mut self) {
338
self.thinking = true;
339
self.thinking_tick = 0;
@@ -444,8 +462,9 @@ mod tui {
462
for msg in &self.messages {
463
lines += wrapped_line_count(&msg.text, msg.role, w);
464
}
447
- if !self.stream_buf.is_empty() {
448
- lines += wrapped_line_count(&self.stream_buf, Role::Assistant, w);
465
+ let visible = self.visible_stream();
466
+ if !visible.is_empty() {
467
+ lines += wrapped_line_count(&visible, Role::Assistant, w);
468
}
469
if self.thinking || self.switching_model {
470
lines += 1;
@@ -851,10 +870,11 @@ mod tui {
870
render_chat_message(&mut lines, msg, inner_width as usize);
871
}
872
854
- if !app.stream_buf.is_empty() {
873
+ let streamed_visible = app.visible_stream();
874
+ if !streamed_visible.is_empty() {
875
let fake = ChatMessage {
876
role: Role::Assistant,
857
- text: app.stream_buf.clone(),
877
+ text: streamed_visible,
878
think_block: None,
879
};
880
render_chat_message(&mut lines, &fake, inner_width as usize);
@@ -1387,7 +1407,36 @@ mod tui {
1407
vec![]
1408
};
1409
1390
- let mut result = match backend.send_message_with_tools(&text, &tools).await {
1410
+ // Bridge the backend's token sink (plain strings) onto the UI update
1411
+ // channel as `Delta` messages. The forwarder lives for the whole turn.
1412
+ let (delta_tx, mut delta_rx) = mpsc::unbounded_channel::<String>();
1413
+ let forward_tx = tx.clone();
1414
+ let forwarder = tokio::spawn(async move {
1415
+ while let Some(piece) = delta_rx.recv().await {
1416
+ if forward_tx
1417
+ .send(InferenceUpdate::Delta(piece))
1418
+ .await
1419
+ .is_err()
1420
+ {
1421
+ break;
1422
+ }
1423
+ }
1424
+ });
1425
+
1426
+ // The first round offers tools, so on-device inference can't stream it
1427
+ // (it must buffer to detect tool calls). With tools disabled there are
1428
+ // none to offer, so it streams directly.
1429
+ let first_sink = if tools.is_empty() {
1430
+ Some(&delta_tx)
1431
+ } else {
1432
+ None
1433
+ };
1434
+ let mut streamed = first_sink.is_some();
1435
+
1436
+ let mut result = match backend
1437
+ .send_message_with_tools(&text, &tools, first_sink)
1438
+ .await
1439
+ {
1440
Ok(r) => r,
1441
Err(err) => {
1442
let _ = tx.send(InferenceUpdate::Error(err)).await;
@@ -1398,6 +1447,8 @@ mod tui {
1447
let mut round = 0;
1448
1449
while !result.tool_calls.is_empty() && round < MAX_TOOL_ROUNDS {
1450
+ // any tool call means the first round didn't produce a final answer
1451
+ streamed = false;
1452
round += 1;
1453
log::info!("tool round {} — {} call(s)", round, result.tool_calls.len());
1454
@@ -1421,14 +1472,24 @@ mod tui {
1472
});
1473
}
1474
1424
- // on the last round, pass no tools so the model must produce text
1475
+ // on the last round, pass no tools so the model must produce text —
1476
+ // that's also the round we can stream on-device.
1477
let next_tools = if round < MAX_TOOL_ROUNDS {
1478
Some(tools.as_slice())
1479
} else {
1480
None
1481
};
1482
+ let sink = if next_tools.is_none() {
1483
+ streamed = true;
1484
+ Some(&delta_tx)
1485
+ } else {
1486
+ None
1487
+ };
1488
1431
- match backend.send_tool_results(tool_results, next_tools).await {
1489
+ match backend
1490
+ .send_tool_results(tool_results, next_tools, sink)
1491
+ .await
1492
+ {
1493
Ok(r) => result = r,
1494
Err(err) => {
1495
let _ = tx.send(InferenceUpdate::Error(err)).await;
@@ -1437,6 +1498,11 @@ mod tui {
1498
}
1499
}
1500
1501
+ // Drop the sink so the forwarder finishes draining any buffered tokens
1502
+ // before we commit the reply.
1503
+ drop(delta_tx);
1504
+ let _ = forwarder.await;
1505
+
1506
if result.tool_calls.is_empty() {
1507
if result.text.is_empty() {
1508
log::warn!(
@@ -1449,6 +1515,9 @@ mod tui {
1515
.to_string(),
1516
))
1517
.await;
1518
+ } else if streamed {
1519
+ // tokens already went out as deltas; just commit the buffer
1520
+ let _ = tx.send(InferenceUpdate::StreamEnd).await;
1521
} else {
1522
let _ = tx.send(InferenceUpdate::Response(result.text)).await;
1523
}
@@ -1588,29 +1657,6 @@ mod tui {
1657
app.tick();
1658
}
1659
1591
- // ── Streaming LLM tokens ──────────────────────────────────────
1592
- chunk = async {
1593
- match app.stream_rx.as_mut() {
1594
- Some(rx) => rx.recv().await,
1595
- None => pending().await,
1596
- }
1597
- } => {
1598
- match chunk {
1599
- Some(chunk) => {
1600
- if !chunk.delta.is_empty() {
1601
- app.push_stream_delta(&chunk.delta);
1602
- }
1603
- if chunk.done {
1604
- app.finalize_stream();
1605
- }
1606
- }
1607
- // sender dropped without done=true
1608
- None => {
1609
- app.finalize_stream();
1610
- }
1611
- }
1612
- }
1613
-
1660
// ── inference updates from background task ───────────────────
1661
update = async {
1662
match app.inference_rx.as_mut() {
@@ -1622,16 +1668,24 @@ mod tui {
1668
Some(InferenceUpdate::ToolUse(name)) => {
1669
app.messages.push(ChatMessage::system(format!("🔧 {name}")));
1670
}
1671
+ Some(InferenceUpdate::Delta(delta)) => {
1672
+ app.push_stream_delta(&delta);
1673
+ }
1674
+ Some(InferenceUpdate::StreamEnd) => {
1675
+ app.finalize_stream();
1676
+ }
1677
Some(InferenceUpdate::Response(text)) => {
1678
app.stop_thinking();
1679
app.messages.push(ChatMessage::assistant(text));
1680
}
1681
Some(InferenceUpdate::Error(msg)) => {
1682
+ app.finalize_stream();
1683
app.stop_thinking();
1684
app.messages.push(ChatMessage::system(format!("error: {msg}")));
1685
}
1686
None => {
1687
// task finished, possibly with no text to show
1688
+ app.finalize_stream();
1689
app.stop_thinking();
1690
}
1691
}
src/main.rs
+126
-29
@@ -61,6 +61,7 @@ use onde::inference::{ChatEngine, GgufModelConfig};
61
62
use crate::backend::{
63
InferenceBackend, LocalBackend, OpenAiBackend, ToolResult as BackendToolResult, ToolSpec,
64
+ TurnResult,
65
};
66
use std::path::PathBuf;
67
use std::sync::atomic::{AtomicBool, Ordering};
@@ -568,6 +569,72 @@ impl SiGitAgent {
569
))
570
}
571
572
+ /// Run one inference turn (`fut`) while concurrently forwarding any streamed
573
+ /// tokens to the editor. The sink receiver is drained as the future runs, so
574
+ /// chunks reach the client live rather than all at once when it resolves.
575
+ ///
576
+ /// `assembled`/`sent`/`streamed_any` persist across the turns of a single
577
+ /// prompt so reasoning is stripped consistently and we never re-send text.
578
+ #[allow(clippy::too_many_arguments)]
579
+ async fn drain_turn<F>(
580
+ &self,
581
+ cx: &ConnectionTo<Client>,
582
+ session_id: &SessionId,
583
+ fut: F,
584
+ sink_rx: &mut tokio::sync::mpsc::UnboundedReceiver<String>,
585
+ assembled: &mut String,
586
+ sent: &mut String,
587
+ streamed_any: &mut bool,
588
+ ) -> Result<TurnResult, backend::BackendError>
589
+ where
590
+ F: std::future::Future<Output = Result<TurnResult, backend::BackendError>>,
591
+ {
592
+ tokio::pin!(fut);
593
+ let result = loop {
594
+ tokio::select! {
595
+ done = &mut fut => break done,
596
+ Some(piece) = sink_rx.recv() => {
597
+ self.emit_visible_chunk(cx, session_id, &piece, assembled, sent, streamed_any);
598
+ }
599
+ }
600
+ };
601
+ // Flush tokens that landed between the last poll and the future resolving.
602
+ while let Ok(piece) = sink_rx.try_recv() {
603
+ self.emit_visible_chunk(cx, session_id, &piece, assembled, sent, streamed_any);
604
+ }
605
+ result
606
+ }
607
+
608
+ /// Append a streamed fragment, strip `<think>` reasoning from the running
609
+ /// text, and send only the newly revealed visible suffix as a chunk. Tracking
610
+ /// the assembled text (not just deltas) keeps think-block stripping correct
611
+ /// even when a tag spans chunk boundaries.
612
+ fn emit_visible_chunk(
613
+ &self,
614
+ cx: &ConnectionTo<Client>,
615
+ session_id: &SessionId,
616
+ piece: &str,
617
+ assembled: &mut String,
618
+ sent: &mut String,
619
+ streamed_any: &mut bool,
620
+ ) {
621
+ assembled.push_str(piece);
622
+ let (_think, visible) = chat::strip_think_blocks(assembled);
623
+ match visible.strip_prefix(sent.as_str()) {
624
+ Some(extra) if !extra.is_empty() => {
625
+ let extra = extra.to_string();
626
+ *sent = visible;
627
+ *streamed_any = true;
628
+ self.send_assistant_message(cx, session_id.clone(), extra)
629
+ .ok();
630
+ }
631
+ // No new visible text, or the visible prefix changed retroactively
632
+ // (rare, e.g. a late-closing think tag): just resync without
633
+ // resending what's already on the wire.
634
+ _ => *sent = visible,
635
+ }
636
+ }
637
+
638
fn send_tool_call_update(
639
&self,
640
cx: &ConnectionTo<Client>,
@@ -1076,8 +1143,26 @@ impl SiGitAgent {
1143
1144
let tools = agent_tools_as_specs();
1145
1079
- let mut result = backend
1080
- .send_message_with_tools(&user_text, &tools)
1146
+ // Token sink: backends stream assistant text through this while a turn
1147
+ // runs. We forward the visible portion to the editor as agent-message
1148
+ // chunks live (see `drain_turn` / `emit_visible_chunk`). The sink stays
1149
+ // alive for the whole prompt so `recv()` only ends when a turn future
1150
+ // resolves, never because every sender was dropped.
1151
+ let (sink, mut sink_rx) = tokio::sync::mpsc::unbounded_channel::<String>();
1152
+ let mut assembled = String::new();
1153
+ let mut sent = String::new();
1154
+ let mut streamed_any = false;
1155
+
1156
+ let mut result = self
1157
+ .drain_turn(
1158
+ cx,
1159
+ &session_id,
1160
+ backend.send_message_with_tools(&user_text, &tools, Some(&sink)),
1161
+ &mut sink_rx,
1162
+ &mut assembled,
1163
+ &mut sent,
1164
+ &mut streamed_any,
1165
+ )
1166
.await
1167
.map_err(|error| {
1168
log::error!("send_message_with_tools failed: {error}");
@@ -1120,39 +1205,51 @@ impl SiGitAgent {
1205
None // last round: force text
1206
};
1207
1123
- result = backend
1124
- .send_tool_results(tool_results, next_tools)
1208
+ result = self
1209
+ .drain_turn(
1210
+ cx,
1211
+ &session_id,
1212
+ backend.send_tool_results(tool_results, next_tools, Some(&sink)),
1213
+ &mut sink_rx,
1214
+ &mut assembled,
1215
+ &mut sent,
1216
+ &mut streamed_any,
1217
+ )
1218
.await
1219
.map_err(|e| agent_client_protocol::Error::new(-32603, e.to_string()))?;
1220
}
1221
1129
- // ── Send the final text response ─────────────────────────────────
1130
- let reply_text = result.text.trim().to_string();
1131
-
1132
- let final_text = if reply_text.is_empty() {
1133
- if round > 0 {
1134
- log::warn!(
1135
- "prompt({}) — model returned empty reply after {} tool round(s)",
1136
- session_id,
1137
- round
1138
- );
1139
- "Something went wrong — the edits didn't go through. Try rephrasing what you need, or point me at the specific lines.".to_string()
1222
+ // ── Final text response ───────────────────────────────────────────
1223
+ // If anything streamed, the visible reply is already on the wire; only
1224
+ // send a trailing block for the non-streamed path (e.g. on-device direct
1225
+ // answers, which onde can't stream while tools are on offer).
1226
+ if !streamed_any {
1227
+ let reply_text = result.text.trim().to_string();
1228
+ let final_text = if reply_text.is_empty() {
1229
+ if round > 0 {
1230
+ log::warn!(
1231
+ "prompt({}) — model returned empty reply after {} tool round(s)",
1232
+ session_id,
1233
+ round
1234
+ );
1235
+ "Something went wrong — the edits didn't go through. Try rephrasing what you need, or point me at the specific lines.".to_string()
1236
+ } else {
1237
+ log::warn!(
1238
+ "prompt({}) — model returned empty reply (no tool rounds)",
1239
+ session_id
1240
+ );
1241
+ String::new()
1242
+ }
1243
} else {
1141
- log::warn!(
1142
- "prompt({}) — model returned empty reply (no tool rounds)",
1143
- session_id
1144
- );
1145
- String::new()
1146
- }
1147
- } else {
1148
- // strip <think> blocks so reasoning tokens stay hidden
1149
- let (_think, visible) = chat::strip_think_blocks(&reply_text);
1150
- visible
1151
- };
1244
+ // strip <think> blocks so reasoning tokens stay hidden
1245
+ let (_think, visible) = chat::strip_think_blocks(&reply_text);
1246
+ visible
1247
+ };
1248
1153
- if !final_text.is_empty() {
1154
- self.send_assistant_message(cx, session_id.clone(), final_text)
1155
- .ok();
1249
+ if !final_text.is_empty() {
1250
+ self.send_assistant_message(cx, session_id.clone(), final_text)
1251
+ .ok();
1252
+ }
1253
}
1254
1255
log::info!("prompt({}) complete — {} tool round(s)", session_id, round);