From 6733720df9b6458ee9508997cb650a06c2832796 Mon Sep 17 00:00:00 2001 From: oxoxDev <164490987+oxoxDev@users.noreply.github.com> Date: Fri, 17 Apr 2026 11:34:00 +0530 Subject: [PATCH] fix(channels): delete thinking messages after final response (#600) (#612) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(channels): add send_channel_delete to REST client (#600) Add DELETE method to BackendOAuthClient for removing channel messages. Used to clean up ephemeral thinking indicators after final response delivery. Co-Authored-By: Claude Opus 4.6 * fix(channels): show thinking ephemerally and delete on final response (#600) Thinking deltas are accumulated and sent as a temporary "💭" message on the channel. When the final agent response is ready the thinking message is deleted via the new DELETE endpoint so the user only sees the clean reply. This replaces the previous behaviour where thinking messages persisted permanently in the Telegram chat. Co-Authored-By: Claude Opus 4.6 * style(channels): apply cargo fmt to bus.rs (#600) Co-Authored-By: Claude Opus 4.6 --------- Co-authored-by: Claude Opus 4.6 --- src/api/rest.rs | 23 ++++ src/openhuman/channels/bus.rs | 203 ++++++++++++++++++++++------------ 2 files changed, 157 insertions(+), 69 deletions(-) diff --git a/src/api/rest.rs b/src/api/rest.rs index 1dffabdf5..ed423271b 100644 --- a/src/api/rest.rs +++ b/src/api/rest.rs @@ -575,6 +575,29 @@ impl BackendOAuthClient { .await } + /// Deletes a message from a communication channel. Used to clean up + /// ephemeral messages (e.g. thinking indicators) after the final + /// response has been delivered. + pub async fn send_channel_delete( + &self, + channel: &str, + message_id: &str, + bearer_jwt: &str, + ) -> Result { + let channel = channel.trim().trim_matches('/'); + anyhow::ensure!(!channel.is_empty(), "channel is required"); + anyhow::ensure!(!message_id.is_empty(), "message_id is required"); + let encoded_channel = urlencoding::encode(channel); + let encoded_id = urlencoding::encode(message_id); + self.authed_json( + bearer_jwt, + Method::DELETE, + &format!("channels/{encoded_channel}/messages/{encoded_id}"), + None, + ) + .await + } + /// Sends a reaction (e.g. emoji) to a message in a channel. pub async fn send_channel_reaction( &self, diff --git a/src/openhuman/channels/bus.rs b/src/openhuman/channels/bus.rs index ca3c98fc4..11d429c27 100644 --- a/src/openhuman/channels/bus.rs +++ b/src/openhuman/channels/bus.rs @@ -141,6 +141,7 @@ impl EventHandler for ChannelInboundSubscriber { "thinking_delta" => { if let Some(delta) = ev.delta.as_ref() { streaming_state.thinking_accumulator.push_str(delta); + streaming_state.thinking_dirty = true; } } "chat_done" | "chat:done" => { @@ -164,20 +165,6 @@ impl EventHandler for ChannelInboundSubscriber { reply_text.len(), streaming_state.message_id, ); - // Send the model's thinking summary as a separate - // message so the user can see the reasoning process. - if !streaming_state.thinking_accumulator.is_empty() { - let summary = format_thinking_summary( - &streaming_state.thinking_accumulator, - ); - tracing::debug!( - "[channel-inbound] sending thinking summary to channel='{}' raw_chars={} summary_chars={}", - channel, - streaming_state.thinking_accumulator.len(), - summary.len(), - ); - send_channel_reply(channel, &summary).await; - } // If we've been streaming progressive edits, replace // the outbound message with the final canonical text. // Otherwise send a fresh message atomically. @@ -211,6 +198,9 @@ impl EventHandler for ChannelInboundSubscriber { } } _ = edit_timer.tick() => { + if streaming_state.thinking_dirty { + flush_thinking_message(channel, &mut streaming_state).await; + } if streaming_state.dirty && !streaming_state.edit_disabled { flush_streaming_edit(channel, &mut streaming_state).await; } @@ -275,9 +265,17 @@ struct StreamingState { /// Latched when the backend doesn't support edits for this channel /// — we stop trying and rely on the final atomic send. edit_disabled: bool, - /// Accumulated model thinking/reasoning text from `thinking_delta` events. - /// Sent as a separate message before the main response at turn completion. + /// Accumulated LLM reasoning from `thinking_delta` events. Shown + /// to the user as an ephemeral "💭 Thinking…" message that is + /// **deleted** once the final response is ready (#600). thinking_accumulator: String, + /// Backend-assigned id of the ephemeral thinking message. Used to + /// delete it at finalization so the user sees only the clean reply. + thinking_message_id: Option, + /// `true` once a thinking message has been posted to the channel. + thinking_sent: bool, + /// New thinking content has arrived since the last thinking flush. + thinking_dirty: bool, } /// Typing-indicator bookkeeping. One per in-flight turn. Latches @@ -332,16 +330,18 @@ async fn send_typing_indicator(channel: &str, state: &mut TypingState) { impl StreamingState { fn compose_draft(&self) -> String { - let mut out = String::new(); - if let Some(ref tool) = self.last_tool { - out.push_str(tool); - out.push('\n'); + let trimmed = self.content.trim_end(); + if trimmed.is_empty() { + // No visible text yet — show a placeholder. Tool indicators + // (🔧 …) are intentionally omitted so the draft only ever + // contains content that is a clean prefix of the final + // response. If the draft persists after finalization the + // user sees benign placeholder text instead of stale tool + // status lines (#600). + "_working…_".to_string() + } else { + trimmed.to_string() } - out.push_str(self.content.trim_end()); - if self.content.is_empty() && self.last_tool.is_none() { - out.push_str("_working…_"); - } - out } } @@ -394,20 +394,6 @@ async fn flush_streaming_edit(channel: &str, state: &mut StreamingState) { } } } else { - // Before posting the first visible draft, deliver any accumulated - // thinking so it appears above the response bubble in the channel. - if !state.thinking_accumulator.is_empty() { - let summary = format_thinking_summary(&state.thinking_accumulator); - tracing::debug!( - "[channel-inbound][stream] sending thinking summary before first draft channel='{}' raw_chars={} summary_chars={}", - channel, - state.thinking_accumulator.len(), - summary.len(), - ); - send_channel_reply(channel, &summary).await; - // Clear so the chat_done handler doesn't send it a second time. - state.thinking_accumulator.clear(); - } let body = json!({ "text": draft }); match client.send_channel_message(channel, &jwt, body).await { Ok(resp) => { @@ -451,6 +437,102 @@ async fn flush_streaming_edit(channel: &str, state: &mut StreamingState) { } } +/// Maximum length of the thinking snippet shown in the ephemeral +/// channel message. Longer reasoning is truncated with "…" to avoid +/// overwhelming the chat. +const MAX_THINKING_DISPLAY_CHARS: usize = 500; + +/// Send or edit the ephemeral "💭 Thinking…" message on the channel. +/// This message is deleted when the final response is ready. +async fn flush_thinking_message(channel: &str, state: &mut StreamingState) { + state.thinking_dirty = false; + + if state.thinking_accumulator.trim().is_empty() { + return; + } + + let mut snippet = state.thinking_accumulator.trim().to_string(); + if snippet.len() > MAX_THINKING_DISPLAY_CHARS { + snippet.truncate(MAX_THINKING_DISPLAY_CHARS); + snippet.push('…'); + } + let text = format!("💭 _{snippet}_"); + + let Some((client, jwt)) = build_channel_client().await else { + return; + }; + + if let Some(ref msg_id) = state.thinking_message_id { + // Edit existing thinking message with updated content. + let body = json!({ "text": text }); + if let Err(err) = client.send_channel_edit(channel, msg_id, &jwt, body).await { + tracing::debug!( + "[channel-inbound][thinking] edit failed channel='{}' msg_id={} err={}", + channel, + msg_id, + err, + ); + } + } else { + // Send initial thinking message. + let body = json!({ "text": text }); + match client.send_channel_message(channel, &jwt, body).await { + Ok(resp) => { + state.thinking_sent = true; + let id = resp + .get("id") + .or_else(|| resp.get("data").and_then(|d| d.get("id"))) + .and_then(|v| v.as_str()) + .map(|s| s.to_string()); + if let Some(id) = id { + tracing::debug!( + "[channel-inbound][thinking] thinking msg sent channel='{}' msg_id={}", + channel, + id, + ); + state.thinking_message_id = Some(id); + } else { + tracing::debug!( + "[channel-inbound][thinking] thinking msg sent but no id returned — will not be deletable", + ); + } + } + Err(err) => { + tracing::warn!( + "[channel-inbound][thinking] failed to send thinking msg channel='{}' err={}", + channel, + err, + ); + } + } + } +} + +/// Delete a previously sent message from the channel. Used to clean +/// up ephemeral thinking messages once the final response is ready. +async fn delete_channel_message(channel: &str, message_id: &str) { + let Some((client, jwt)) = build_channel_client().await else { + return; + }; + match client.send_channel_delete(channel, message_id, &jwt).await { + Ok(_) => { + tracing::info!( + "[channel-inbound] deleted ephemeral msg channel='{}' msg_id={}", + channel, + message_id, + ); + } + Err(err) => { + tracing::warn!( + "[channel-inbound] failed to delete ephemeral msg channel='{}' msg_id={} err={}", + channel, + message_id, + err, + ); + } + } +} + /// Deliver the final canonical reply. /// /// **Invariant**: if a draft message has already been posted to the @@ -461,6 +543,13 @@ async fn flush_streaming_edit(channel: &str, state: &mut StreamingState) { /// creates a fresh outbound message is when no draft has been posted /// at all. async fn finalize_channel_reply(channel: &str, state: &mut StreamingState, final_text: &str) { + // ── Clean up ephemeral thinking message ────────────────────── + // Delete the "💭 Thinking…" message so the user only sees the + // clean final response (#600). + if let Some(ref thinking_id) = state.thinking_message_id { + delete_channel_message(channel, thinking_id).await; + } + if let Some(ref message_id) = state.message_id { // We committed to a draft earlier in the turn. Always attempt // to edit it with the canonical reply, even when we'd @@ -503,12 +592,15 @@ async fn finalize_channel_reply(channel: &str, state: &mut StreamingState, final } if state.draft_sent { // A draft was posted but the backend didn't return an id, so - // we have nothing to edit. Posting a fresh message here would - // give the user two bubbles — skip silently. + // we have nothing to edit. Since the draft only contains a + // clean text prefix (or "_working…_" placeholder), sending the + // final response as a second bubble is acceptable — leaving + // the user without the canonical reply is worse (#600). tracing::warn!( - "[channel-inbound] skipping fresh send on channel='{}' — an id-less draft was already posted earlier this turn (duplicate prevented)", + "[channel-inbound] sending fresh reply on channel='{}' — id-less draft exists but user needs the final response", channel, ); + send_channel_reply(channel, final_text).await; return; } // No draft exists — this is the first (and only) message for the @@ -598,33 +690,6 @@ async fn send_channel_reply(channel: &str, text: &str) { } } -/// Maximum characters of raw thinking content included in the summary sent -/// to the channel. Telegram's hard message limit is 4 096 chars; keeping the -/// body well below that leaves room for the header and the trailing ellipsis. -const MAX_THINKING_CHARS: usize = 1500; - -/// Format accumulated thinking content into a readable message for delivery -/// to the channel. Truncates at the last word boundary when the raw content -/// exceeds [`MAX_THINKING_CHARS`], appending `…` to signal the cut. -fn format_thinking_summary(thinking: &str) -> String { - let trimmed = thinking.trim(); - let body = if trimmed.len() > MAX_THINKING_CHARS { - // Find a safe char boundary at or before MAX_THINKING_CHARS so we never - // slice in the middle of a multi-byte UTF-8 sequence. - let safe_end = (0..=MAX_THINKING_CHARS) - .rev() - .find(|&i| trimmed.is_char_boundary(i)) - .unwrap_or(0); - let slice = &trimmed[..safe_end]; - // Back off further to the last whitespace to avoid cutting mid-word. - let boundary = slice.rfind(|c: char| c.is_whitespace()).unwrap_or(safe_end); - format!("{}…", &slice[..boundary]) - } else { - trimmed.to_string() - }; - format!("💭 Thinking:\n\n{}", body) -} - #[cfg(test)] mod tests { use super::*;