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
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
//! The streaming edit loop, extracted from `handler.rs` (#1086 seam 5).
//!
//! Spawned once per turn: every 1.5s it snapshots the shared `StreamingState`
//! under lock, folds queued display items (tools, intermediates) into the
//! chat in chronological order, re-sticks the open flow block when newer
//! chatter buries it (#451), refreshes tool-group messages whose status
//! changed, and edits the response message in place. It exits when the turn
//! cancels the token.
use std::sync::Arc;
use teloxide::prelude::*;
use teloxide::types::{ChatAction, ChatId, MessageId, ParseMode, ThreadId};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use super::flow::{
DisplayItem, StreamingState, append_intermediate_to_flow, append_tool_group,
restick_flow_if_buried,
};
use super::handler::{fire_reaction, thinking_status_excerpt};
use super::markdown::markdown_to_telegram_html;
use super::send::{best_effort_delete, fire_chat_action, message_in_thread};
use super::state::TelegramState;
use crate::brain::AgentService;
#[allow(clippy::too_many_arguments)]
pub(crate) fn spawn_edit_loop(
bot: &Bot,
chat: ChatId,
react_target: Option<MessageId>,
thread_id: Option<ThreadId>,
is_dm: bool,
streaming: &Arc<std::sync::Mutex<StreamingState>>,
edit_cancel: &CancellationToken,
telegram_state: &Arc<TelegramState>,
agent: &Arc<AgentService>,
session_id: Uuid,
) -> tokio::task::JoinHandle<()> {
tokio::spawn({
let bot = bot.clone();
let st = streaming.clone();
let cancel = edit_cancel.clone();
let tg = telegram_state.clone();
let agent = agent.clone();
let sid = session_id;
async move {
loop {
tokio::select! {
_ = cancel.cancelled() => break,
_ = tokio::time::sleep(std::time::Duration::from_millis(1500)) => {
// ── Snapshot state under lock, then release immediately ──
struct Snapshot {
dirty: bool,
recreate: bool,
response_text: String,
msg_id: Option<MessageId>,
tool_round_count: usize,
/// Ordered display items (tools + intermediates in chronological order)
display_items: Vec<DisplayItem>,
/// Dirty tools that already have messages (need editing, not new sends)
tool_edits: Vec<(usize, String, Option<bool>, MessageId)>,
has_active_tools: bool,
processing: bool,
/// Short excerpt of the latest reasoning chunk used as
/// a context-aware status line during the pre-tool
/// phase. Falls back to a fun-quip rotation when
/// reasoning hasn't started yet.
thinking_excerpt: Option<String>,
}
let mut settle_flow = false;
let snap = {
let mut s = st.lock().unwrap_or_else(|e| e.into_inner());
let has_display = !s.display_queue.is_empty();
let any_tools_dirty = s.tool_msgs.iter().any(|t| t.dirty);
let has_active_tools = s.tool_msgs.iter().any(|t| t.completed.is_none());
let processing = s.processing;
if !s.dirty && !s.recreate && !any_tools_dirty && !has_display && !has_active_tools && !processing { continue; }
// Drain the ordered display queue
let display_items: Vec<DisplayItem> = s.display_queue.drain(..).collect();
// Collect dirty tools that already have messages (for editing)
let tool_edits: Vec<_> = s.tool_msgs.iter().enumerate()
.filter(|(_, t)| t.dirty && t.msg_id.is_some())
.map(|(i, t)| {
let label = format!("**{}**{}", t.name, t.context);
(i, label, t.completed, t.msg_id.unwrap())
})
.collect();
// Mark tools as not dirty
for t in s.tool_msgs.iter_mut().filter(|t| t.dirty) {
t.dirty = false;
}
// Snapshot response
let response_text = if s.dirty || s.recreate {
s.render()
} else {
String::new()
};
let snap = Snapshot {
dirty: s.dirty,
recreate: s.recreate,
response_text,
msg_id: s.msg_id,
tool_round_count: s.tool_round_count,
display_items,
tool_edits,
has_active_tools,
processing,
thinking_excerpt: thinking_status_excerpt(&s.thinking),
};
// Pre-clear state that will be handled
if s.recreate {
s.recreate = false;
}
if s.dirty {
s.dirty = false;
}
// Clear status tracking only when final response arrives (#313)
// Don't clear on intermediates — keep the status message alive and
// edit it in place throughout multi-tool sequences, so we get one
// updating message instead of N+1 separate messages.
if snap.dirty && !snap.response_text.is_empty() {
s.tools_started_at = None;
s.tool_round_count = 0;
// Header settles to the plain "N tool calls"
// via an immediate refresh below (#360).
if s.flow_status.take().is_some() && s.open_group_msg_id.is_some()
{
settle_flow = true;
}
}
snap
};
// Lock is now released
// ── Ordered display: tools and intermediates in chronological order ──
// Buffer consecutive tool calls to group them into collapsible blocks
let mut tool_buffer: Vec<usize> = Vec::new();
for item in &snap.display_items {
match item {
DisplayItem::NewTool(idx) => {
// Buffer this tool call
tool_buffer.push(*idx);
}
DisplayItem::Intermediate(text) => {
// Flush buffered tools into the open flow,
// then fold this intermediate into the SAME
// in-place processing-log message. It no
// longer lands as its own message, so only
// the final response stays clean at the
// bottom (#300).
append_tool_group(&bot, chat, thread_id, &st, &tool_buffer)
.await;
tool_buffer.clear();
// Sanitize exactly as before folding:
// strip LLM artifacts, redact secrets, strip
// <<IMG:>> markers (the final-response
// handler sends the image), and extract +
// fire <<react:>> now so a mid-turn reaction
// acknowledges the user immediately (#261).
let text = crate::utils::sanitize::strip_llm_artifacts(text);
let text = crate::utils::redact_secrets_scoped(&text, is_dm);
let (text, _img_paths) =
crate::utils::extract_img_markers(&text);
let (text, react_emoji) =
crate::utils::extract_react_marker(&text);
// A resumed turn has no inbound message to
// react to: the marker is stripped above but
// nothing fires (#261).
if let Some(ref emoji) = react_emoji
&& let Some(target) = react_target
{
fire_reaction(&bot, chat, target, emoji).await;
}
// A substantial rich report (a table) the
// model emits before a tool call would be
// buried in the collapsed log — surface it as
// its own rich message instead (#582). Thin
// narration keeps folding: folded intermediates
// are NOT recorded in sent_intermediates, so
// the final-response dedup does not suppress the
// visible answer just because it also appears in
// the collapsed trace.
if super::intermediates::is_deliverable_rich_report(&text) {
super::intermediates::deliver_intermediate_message(
&bot, chat, thread_id, &st, &text,
)
.await;
} else {
append_intermediate_to_flow(
&bot, chat, thread_id, &st, &text,
)
.await;
}
}
}
}
// Flush any remaining buffered tools into the open group.
// No close here: the run may continue on the next tick, in
// which case those tools append to this same message.
append_tool_group(&bot, chat, thread_id, &st, &tool_buffer).await;
// ── Re-stick the open block to the bottom if buried (#451) ──
// A new round landed this tick (tools/intermediates were in
// the display queue). If newer chatter has pushed the block
// above the newest message, relocate it to the bottom. Gated
// on real appends, never plain status ticks, so an idle chat
// sees no churn.
if !snap.display_items.is_empty() {
let newest = tg.newest_incoming_msg_id(chat.0);
restick_flow_if_buried(&bot, chat, thread_id, &st, newest).await;
}
// ── Update tool-group messages for tools that changed status ──
// A completed tool shares its group's message with its
// siblings, so re-render the whole group (never a single
// tool line, which would overwrite the block). Refresh each
// distinct group once.
// A tool status flip (⚙️ → ✅/❌) re-renders the whole
// processing-log flow (tools + folded intermediates) in
// its single message.
// Show progress when: tools are active, OR tools ran but no
// response yet, OR still processing (initial wait).
let show_status = snap.has_active_tools
|| (snap.tool_round_count > 0 && snap.response_text.is_empty())
|| snap.processing;
// ── Single progress surface: the flow message ──
// The live status (thinking / Working-on / activity
// preview), wall-clock duration, and plan/goal/ctx
// sections all ride the flow header (#360, #480,
// #509). While no flow is open and the turn is still
// working, the shared tick opens it header-only on
// this activity tick; the legacy pre-block status
// bubble is gone.
let turn_done = snap.dirty && !snap.response_text.is_empty();
// Only show the thinking excerpt as a status preview.
// The user's message is NOT what the bot is "working on"
// it's just the input request, so showing it as
// "Working on: <user message>" is confusing. The goal
// section (from GoalManager) already shows what the bot
// is actually working on when a plan task is active.
let preview = snap
.thinking_excerpt
.as_deref()
.map(|t| format!("🧠 {t}"));
let flow_needs_refresh = !snap.tool_edits.is_empty() || settle_flow;
super::flow_chrome::tick_flow_header(
&bot,
chat,
thread_id,
&st,
&agent,
sid,
show_status,
turn_done,
preview,
flow_needs_refresh,
)
.await;
// Update the persistent plan card in place (#580): the
// checklist lives on its own card now, not the flow
// block, so it advances here as tasks complete.
let plan_kb = {
st.lock().unwrap_or_else(|e| e.into_inner()).sections.plan_kb
};
super::plan_card::refresh_plan_card(
&bot, chat, thread_id, &tg, &agent, sid, plan_kb,
)
.await;
// ── Response message (thinking + response, always at bottom) ──
// Stale-placeholder cleanup runs unconditionally: a bubble
// opened before the first tool call must still be removed
// once a block opens.
if snap.recreate
&& let Some(old_mid) = snap.msg_id
{
best_effort_delete(&bot, chat, old_mid, "recreate swap").await;
let mut s = st.lock().unwrap_or_else(|e| e.into_inner());
s.msg_id = None;
}
// While a processing-log block is open, mid-round narration
// folds into that block (append_intermediate_to_flow) and the
// final answer is delivered by deliver_final_response at turn
// end. Opening a standalone streaming bubble here leaks the
// intermediate text as its own message beneath the folded
// block (#490), so only stream the placeholder when NO
// processing-log block is open. Re-read the id: the
// header tick above may have just opened the flow.
let open_block = {
let s = st.lock().unwrap_or_else(|e| e.into_inner());
s.open_group_msg_id
};
if (snap.dirty || snap.recreate)
&& open_block.is_none()
&& !snap.response_text.is_empty()
{
let current_msg_id = {
let s = st.lock().unwrap_or_else(|e| e.into_inner());
s.msg_id
};
if current_msg_id.is_none() {
// Success-silent until #1085: the twin in
// resume.rs logged both outcomes while this
// one dropped the error, so a failing
// placeholder send was invisible on the
// handler path.
match message_in_thread(&bot, chat, thread_id, "\u{258b}").await {
Ok(m) => {
super::telemetry::log_send_success(
"turn",
"-",
"-",
"placeholder",
"new",
chat.0,
thread_id.map(|t| t.0.0),
m.id.0,
"\u{258b}".len(),
&super::telemetry::content_hash8("\u{258b}"),
);
let mut s = st.lock().unwrap_or_else(|e| e.into_inner());
s.msg_id = Some(m.id);
}
Err(e) => {
super::telemetry::log_send_failure(
"turn",
"-",
"-",
"placeholder",
"new",
chat.0,
thread_id.map(|t| t.0.0),
"\u{258b}".len(),
&super::telemetry::content_hash8("\u{258b}"),
&e.to_string(),
);
}
}
}
let msg_id = {
let s = st.lock().unwrap_or_else(|e| e.into_inner());
s.msg_id
};
if let Some(mid) = msg_id {
// Strip any complete <<react:emoji>>
// directive from the streaming snapshot so
// the raw marker never flashes in the
// placeholder (#261). The reaction itself
// fires from the intermediate/final paths.
let (clean, _) =
crate::utils::extract_react_marker(&snap.response_text);
let html = markdown_to_telegram_html(&clean);
let display = format!("{}\u{258b}", html); // ▋ cursor
if let Err(e) = bot
.edit_message_text(chat, mid, display)
.parse_mode(ParseMode::Html)
.await
{
// Review F10: placeholder edits were fully
// silent; a failing edit stream is now visible.
tracing::warn!(
"Telegram: streaming placeholder edit failed (chat={} msg={}): {}",
chat.0,
mid.0,
e
);
}
}
}
// Re-send typing indicator after any bot message
fire_chat_action(&bot, chat, thread_id, ChatAction::Typing, "post-message typing").await;
}
}
}
}
})
}