From a20ebcca0af2c01d15e7f70dfbfb9c39a90b13b0 Mon Sep 17 00:00:00 2001 From: homelab Date: Thu, 17 Sep 2026 14:56:18 -0600 Subject: [PATCH 01/10] feat(telegram): anchor replies to the trigger message --- src/audit.rs | 1 + src/channel.rs | 13 +++++++++++-- src/gateway/mod.rs | 13 ++++++++++--- src/gateway/tests.rs | 4 ++++ src/gateway/worker.rs | 3 ++- src/slack.rs | 1 + src/telegram.rs | 28 ++++++++++++++++++++++++++++ 7 files changed, 57 insertions(+), 6 deletions(-) diff --git a/src/audit.rs b/src/audit.rs index eb87454..5c01aa8 100644 --- a/src/audit.rs +++ b/src/audit.rs @@ -395,6 +395,7 @@ mod tests { images: Vec::new(), is_from_me: false, is_supported: true, + reply_to_message_id: None, thread_id: None, } } diff --git a/src/channel.rs b/src/channel.rs index 4dc8720..9fb098b 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -55,6 +55,9 @@ pub struct RawMessage { pub is_supported: bool, /// Channel-specific thread/topic id (Telegram `message_thread_id`). pub thread_id: Option, + /// Provider message id of the inbound message (Telegram `message_id`), + /// used to anchor replies via `reply_parameters`. + pub reply_to_message_id: Option, } impl RawMessage { @@ -70,6 +73,8 @@ impl RawMessage { pub struct OutboundChunk { pub text: String, pub rich_markdown: bool, + /// Telegram message id to reply to (anchoring), when known. + pub reply_to_message_id: Option, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -422,6 +427,7 @@ impl ChannelContract for IMessageChannel { .collect(), is_from_me: message.is_from_me, is_supported: true, + reply_to_message_id: None, thread_id: None, }) .collect()) @@ -638,9 +644,9 @@ impl ChannelContract for Telegram { async fn send_chunk(&self, target: &str, chunk: &OutboundChunk) -> Result<()> { if chunk.rich_markdown { - self.send_rich(target, &chunk.text).await + self.send_rich_reply(target, &chunk.text, chunk.reply_to_message_id).await } else { - self.send_plain(target, &chunk.text).await + self.send_plain_reply(target, &chunk.text, chunk.reply_to_message_id).await } } @@ -883,6 +889,7 @@ mod tests { images: Vec::new(), is_from_me, is_supported: true, + reply_to_message_id: None, thread_id: None, } } @@ -900,6 +907,7 @@ mod tests { images: Vec::new(), is_from_me: false, is_supported: true, + reply_to_message_id: None, thread_id: None, } } @@ -995,6 +1003,7 @@ mod tests { images: Vec::new(), is_from_me: false, is_supported: true, + reply_to_message_id: None, thread_id: None, }; diff --git a/src/gateway/mod.rs b/src/gateway/mod.rs index e02b687..270a71f 100644 --- a/src/gateway/mod.rs +++ b/src/gateway/mod.rs @@ -44,6 +44,8 @@ struct Job { voice_attachment: Option, image_attachments: Vec, approval_origin: AnswerOrigin, + /// Telegram message id of the trigger message (reply anchoring). + telegram_reply_anchor: Option, } /// Shared, cheaply cloneable context handed to each worker task. @@ -61,6 +63,8 @@ struct Ctx { audit: Arc, voice: Option, schedule_destination: Option, + /// Telegram message id to anchor replies to, set by the worker per turn. + telegram_reply_anchor: Arc>>, #[cfg(test)] setup_failure_replies: Arc>>, #[cfg(test)] @@ -412,6 +416,7 @@ impl Gateway { channel: channel.clone(), run_timeout: cfg.run_timeout_dur()?, reply_marker: crate::channel::REPLY_MARKER.to_string(), + telegram_reply_anchor: Arc::new(Mutex::new(None)), assistant_dir: cfg.assistant_dir.clone(), audit, schedule_destination: None, @@ -900,6 +905,7 @@ impl Gateway { voice_attachment: m.voice.clone(), image_attachments: m.images.clone(), approval_origin, + telegram_reply_anchor: m.reply_to_message_id, }; if job.image_attachments.is_empty() && job.text.trim().eq_ignore_ascii_case("/stop") { @@ -1231,7 +1237,7 @@ async fn reply_to(ctx: &Ctx, target: &str, text: &str) -> bool { return false; } for chunk in chunks { - if let Err(error) = send_reply_chunk(ctx, target, &chunk).await { + if let Err(error) = send_reply_chunk(ctx, target, &mut chunk.clone()).await { error!("send error to {target}: {error}"); return false; } @@ -1242,7 +1248,7 @@ async fn reply_to(ctx: &Ctx, target: &str, text: &str) -> bool { async fn send_reply_chunk( ctx: &Ctx, target: &str, - chunk: &crate::channel::OutboundChunk, + chunk: &mut crate::channel::OutboundChunk, ) -> Result<()> { #[cfg(test)] { @@ -1277,6 +1283,7 @@ async fn send_reply_chunk( } #[cfg(not(test))] { + chunk.reply_to_message_id = *ctx.telegram_reply_anchor.lock().unwrap(); let timeout = ctx.channel.delivery_semantics().send_timeout; if timeout.is_zero() { ctx.channel.send_chunk(target, chunk).await @@ -1325,7 +1332,7 @@ async fn send_scheduled_chunk( target: &str, chunk: &crate::channel::OutboundChunk, ) -> Result<()> { - send_reply_chunk(ctx, target, chunk).await + send_reply_chunk(ctx, target, &mut chunk.clone()).await } fn complete_row(store: &Arc>, ack: &Arc>, channel: &str, row_id: i64) { diff --git a/src/gateway/tests.rs b/src/gateway/tests.rs index c7a0bd0..67b6911 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -166,6 +166,7 @@ fn msg(chat: &str, handle: &str, from_me: bool, text: &str) -> RawMessage { is_from_me: from_me, is_group: false, is_supported: true, + reply_to_message_id: None, thread_id: None, } } @@ -4096,6 +4097,7 @@ fn message(row_id: i64, chat: &str, handle: &str, is_from_me: bool, text: &str) is_from_me, is_group: false, is_supported: true, + reply_to_message_id: None, thread_id: None, } } @@ -4119,6 +4121,7 @@ fn telegram_message( is_from_me: false, is_group, is_supported: true, + reply_to_message_id: None, thread_id: None, } } @@ -4173,6 +4176,7 @@ fn slack_image_message( is_from_me: false, is_group: false, is_supported: true, + reply_to_message_id: None, thread_id: None, } } diff --git a/src/gateway/worker.rs b/src/gateway/worker.rs index ecba8df..28bb16e 100644 --- a/src/gateway/worker.rs +++ b/src/gateway/worker.rs @@ -848,6 +848,7 @@ pub(super) async fn record_and_deliver( origin: OutboundOrigin, text: &str, ) -> Result { + *ctx.telegram_reply_anchor.lock().unwrap() = job.telegram_reply_anchor; let outbound = ctx.history.lock().unwrap().record_outbound( job.inbound_id, origin, @@ -962,7 +963,7 @@ async fn deliver_outbound_once( .enumerate() .skip(outbound.delivery_chunk_index) { - if let Err(error) = super::send_reply_chunk(ctx, target, chunk).await { + if let Err(error) = super::send_reply_chunk(ctx, target, &mut chunk.clone()).await { error!( "outbound {} chunk {index} send error to {target}: {error}", outbound.id diff --git a/src/slack.rs b/src/slack.rs index 8e47080..7a2296c 100644 --- a/src/slack.rs +++ b/src/slack.rs @@ -663,6 +663,7 @@ impl Inbox { .collect(), is_from_me: row.get(8)?, is_supported: row.get(9)?, + reply_to_message_id: None, thread_id: None, }) })? diff --git a/src/telegram.rs b/src/telegram.rs index 2505e34..9f4d703 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -263,6 +263,15 @@ impl Telegram { } pub async fn send_rich(&self, target: &str, text: &str) -> Result<()> { + self.send_rich_reply(target, text, None).await + } + + pub async fn send_rich_reply( + &self, + target: &str, + text: &str, + reply_to_message_id: Option, + ) -> Result<()> { if text.encode_utf16().count() > TEXT_LIMIT { bail!("Telegram rich message exceeds the {TEXT_LIMIT} character chunk limit"); } @@ -270,6 +279,9 @@ impl Telegram { let mut payload = target_payload(target); payload["text"] = json!(html); payload["parse_mode"] = json!("HTML"); + if let Some(id) = reply_to_message_id { + payload["reply_parameters"] = json!({"message_id": id, "allow_sending_without_reply": true}); + } let transport_response = self .post_with_topic_fallback("sendMessage", payload) .await?; @@ -285,8 +297,20 @@ impl Telegram { } pub async fn send_plain(&self, target: &str, text: &str) -> Result<()> { + self.send_plain_reply(target, text, None).await + } + + pub async fn send_plain_reply( + &self, + target: &str, + text: &str, + reply_to_message_id: Option, + ) -> Result<()> { let mut payload = target_payload(target); payload["text"] = json!(text); + if let Some(id) = reply_to_message_id { + payload["reply_parameters"] = json!({"message_id": id, "allow_sending_without_reply": true}); + } let transport_response = self .post_with_topic_fallback("sendMessage", payload) .await?; @@ -517,6 +541,7 @@ impl Update { is_from_me: false, is_supported: false, thread_id: None, + reply_to_message_id: None, }; }; let images = message @@ -572,6 +597,7 @@ impl Update { is_from_me: false, is_supported: true, thread_id: message.message_thread_id, + reply_to_message_id: message.message_id, } } } @@ -582,6 +608,8 @@ struct TelegramMessage { from: Option, chat: Chat, #[serde(default)] + message_id: Option, + #[serde(default)] text: Option, #[serde(default)] caption: Option, From 7862433c235db8746cc80d5d217dc5c470a92730 Mon Sep 17 00:00:00 2001 From: homelab Date: Thu, 17 Sep 2026 15:04:59 -0600 Subject: [PATCH 02/10] test: add message_id to TelegramMessage test initializers --- src/telegram.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/telegram.rs b/src/telegram.rs index 9f4d703..738f737 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -895,6 +895,7 @@ mod tests { document: None, voice: None, message_thread_id: None, + message_id: 1, }), } .into_raw(); @@ -926,6 +927,7 @@ mod tests { mime_type: Some("audio/ogg".to_string()), }), message_thread_id: None, + message_id: 1, }), } .into_raw(); From d207a500e78d195e677c3c6ed13cc55ab27e8084 Mon Sep 17 00:00:00 2001 From: homelab Date: Thu, 17 Sep 2026 15:09:38 -0600 Subject: [PATCH 03/10] test: update initializers for telegram reply-anchor fields --- src/channel.rs | 5 +++++ src/gateway/tests.rs | 4 ++++ src/telegram.rs | 4 ++-- 3 files changed, 11 insertions(+), 2 deletions(-) diff --git a/src/channel.rs b/src/channel.rs index 9fb098b..0eb4b9e 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -521,6 +521,7 @@ impl ChannelContract for IMessageChannel { Vec::new() } else { vec![OutboundChunk { + reply_to_message_id: None, text: format!("{text}{marker}"), rich_markdown: false, }] @@ -636,6 +637,7 @@ impl ChannelContract for Telegram { crate::telegram::split_text(text) .into_iter() .map(|text| OutboundChunk { + reply_to_message_id: None, text, rich_markdown: true, }) @@ -742,6 +744,7 @@ impl ChannelContract for Slack { crate::slack::split_text(&crate::markdown::to_slack_mrkdwn_for_chunking(text)) .into_iter() .map(|text| OutboundChunk { + reply_to_message_id: None, text, rich_markdown: true, }) @@ -947,6 +950,7 @@ mod tests { assert_eq!( channel.outbound_chunks("reply", REPLY_MARKER), [OutboundChunk { + reply_to_message_id: None, text: format!("reply{REPLY_MARKER}"), rich_markdown: false, }] @@ -982,6 +986,7 @@ mod tests { assert_eq!( channel.outbound_chunks("reply", REPLY_MARKER), [OutboundChunk { + reply_to_message_id: None, text: "reply".to_string(), rich_markdown: true, }] diff --git a/src/gateway/tests.rs b/src/gateway/tests.rs index 67b6911..1475abd 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -208,6 +208,7 @@ fn setup_failure_ctx( .unwrap(); assert_eq!(inbound_id, 1); Ctx { + telegram_reply_anchor: std::sync::Arc::new(std::sync::Mutex::new(None)), cfg: test_config( &temp_path("setup-failure-state").to_string_lossy(), &temp_path("setup-failure-sessions").to_string_lossy(), @@ -240,6 +241,7 @@ fn setup_failure_ctx( fn setup_failure_job(row_id: i64) -> Job { Job { + telegram_reply_anchor: None, row_id, inbound_id: 1, thread: "imessage:self:me".to_string(), @@ -2787,6 +2789,7 @@ async fn closed_worker_queue_is_recovered_without_another_message() { .unwrap(); gateway.ack.lock().unwrap().in_flight.insert(1); let lost_job = Job { + telegram_reply_anchor: None, row_id: 1, inbound_id: lost_inbound_id, thread: thread.to_string(), @@ -3097,6 +3100,7 @@ async fn stop_targets_the_current_row_ahead_of_retained_failures() { .unwrap(); let thread = "imessage:self:me@icloud.com"; let make_job = |row_id, inbound_id, text: &str| Job { + telegram_reply_anchor: None, row_id, inbound_id, thread: thread.to_string(), diff --git a/src/telegram.rs b/src/telegram.rs index 738f737..b0fe522 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -895,7 +895,7 @@ mod tests { document: None, voice: None, message_thread_id: None, - message_id: 1, + message_id: Some(1), }), } .into_raw(); @@ -927,7 +927,7 @@ mod tests { mime_type: Some("audio/ogg".to_string()), }), message_thread_id: None, - message_id: 1, + message_id: Some(1), }), } .into_raw(); From 5d1fc68b3772eb6c8772b24ec9418efc89d0a5a6 Mon Sep 17 00:00:00 2001 From: Bernardo Kuri Date: Thu, 17 Sep 2026 15:30:21 -0600 Subject: [PATCH 04/10] fix(gateway): pass telegram reply anchor explicitly, never anchor scheduled sends The reply anchor rode in a shared Arc>> slot on Ctx: concurrent worker turns could overwrite each other's anchor mid-delivery, and scheduled/proactive sends inherited whatever anchor the last interactive job left behind. Thread the anchor as a parameter from the job through the delivery chain instead. Scheduled sends pass None and are never anchored; approval confirmations anchor to the message they answer; schedule questions, primary delivery, and other system replies are unanchored. Adds a payload test asserting reply_parameters is present only when anchored, with allow_sending_without_reply set. --- src/channel.rs | 6 ++++-- src/gateway/mod.rs | 29 ++++++++++++++++------------- src/gateway/tests.rs | 3 +-- src/gateway/worker.rs | 20 ++++++++++++++------ src/telegram.rs | 29 +++++++++++++++++++++++++++-- 5 files changed, 62 insertions(+), 25 deletions(-) diff --git a/src/channel.rs b/src/channel.rs index 0eb4b9e..80a0985 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -646,9 +646,11 @@ impl ChannelContract for Telegram { async fn send_chunk(&self, target: &str, chunk: &OutboundChunk) -> Result<()> { if chunk.rich_markdown { - self.send_rich_reply(target, &chunk.text, chunk.reply_to_message_id).await + self.send_rich_reply(target, &chunk.text, chunk.reply_to_message_id) + .await } else { - self.send_plain_reply(target, &chunk.text, chunk.reply_to_message_id).await + self.send_plain_reply(target, &chunk.text, chunk.reply_to_message_id) + .await } } diff --git a/src/gateway/mod.rs b/src/gateway/mod.rs index 270a71f..2484725 100644 --- a/src/gateway/mod.rs +++ b/src/gateway/mod.rs @@ -63,8 +63,6 @@ struct Ctx { audit: Arc, voice: Option, schedule_destination: Option, - /// Telegram message id to anchor replies to, set by the worker per turn. - telegram_reply_anchor: Arc>>, #[cfg(test)] setup_failure_replies: Arc>>, #[cfg(test)] @@ -212,7 +210,7 @@ impl GatewayGroup { .iter() .find(|gateway| gateway.channel.id() == destination.channel) .context("resolved primary delivery channel is unavailable")?; - if !reply_to(&gateway.ctx, &destination.target, text).await { + if !reply_to(&gateway.ctx, &destination.target, text, None).await { anyhow::bail!( "primary delivery to {} target {:?} failed", destination.channel, @@ -416,7 +414,6 @@ impl Gateway { channel: channel.clone(), run_timeout: cfg.run_timeout_dur()?, reply_marker: crate::channel::REPLY_MARKER.to_string(), - telegram_reply_anchor: Arc::new(Mutex::new(None)), assistant_dir: cfg.assistant_dir.clone(), audit, schedule_destination: None, @@ -460,7 +457,7 @@ impl Gateway { .lock() .unwrap() .create_question(&question, now_ms())?; - let delivered = reply_to(&self.ctx, &question.target, &question.render_text()).await; + let delivered = reply_to(&self.ctx, &question.target, &question.render_text(), None).await; self.ctx.history.lock().unwrap().mark_question_delivery( &id, if delivered { @@ -792,7 +789,9 @@ impl Gateway { | jobs::ScheduleDecision::NotScheduleReview => None, }; if let Some(confirmation) = confirmation { - if !reply_to(&self.ctx, &target, &confirmation).await { + if !reply_to(&self.ctx, &target, &confirmation, m.reply_to_message_id) + .await + { warn!("[{thread}] schedule review confirmation delivery failed"); } } @@ -829,6 +828,7 @@ impl Gateway { &self.ctx, &target, "Job approval is no longer used. Ask me to create the job again and I will write it directly to the assistant repository.", + m.reply_to_message_id, ) .await { @@ -1230,14 +1230,14 @@ fn audit_schedule_events(ctx: &Ctx, ledger: &mut jobs::Ledger) { } } -async fn reply_to(ctx: &Ctx, target: &str, text: &str) -> bool { +async fn reply_to(ctx: &Ctx, target: &str, text: &str, reply_anchor: Option) -> bool { let chunks = ctx.channel.outbound_chunks(text, &ctx.reply_marker); if chunks.is_empty() { error!("send error to {target}: channel produced no outbound chunks"); return false; } for chunk in chunks { - if let Err(error) = send_reply_chunk(ctx, target, &mut chunk.clone()).await { + if let Err(error) = send_reply_chunk(ctx, target, reply_anchor, &chunk).await { error!("send error to {target}: {error}"); return false; } @@ -1245,10 +1245,12 @@ async fn reply_to(ctx: &Ctx, target: &str, text: &str) -> bool { true } +#[cfg_attr(test, allow(unused_variables))] async fn send_reply_chunk( ctx: &Ctx, target: &str, - chunk: &mut crate::channel::OutboundChunk, + reply_anchor: Option, + chunk: &crate::channel::OutboundChunk, ) -> Result<()> { #[cfg(test)] { @@ -1283,12 +1285,13 @@ async fn send_reply_chunk( } #[cfg(not(test))] { - chunk.reply_to_message_id = *ctx.telegram_reply_anchor.lock().unwrap(); + let mut chunk = chunk.clone(); + chunk.reply_to_message_id = reply_anchor; let timeout = ctx.channel.delivery_semantics().send_timeout; if timeout.is_zero() { - ctx.channel.send_chunk(target, chunk).await + ctx.channel.send_chunk(target, &chunk).await } else { - match tokio::time::timeout(timeout, ctx.channel.send_chunk(target, chunk)).await { + match tokio::time::timeout(timeout, ctx.channel.send_chunk(target, &chunk)).await { Ok(result) => result, Err(_) => anyhow::bail!("send timed out"), } @@ -1332,7 +1335,7 @@ async fn send_scheduled_chunk( target: &str, chunk: &crate::channel::OutboundChunk, ) -> Result<()> { - send_reply_chunk(ctx, target, &mut chunk.clone()).await + send_reply_chunk(ctx, target, None, chunk).await } fn complete_row(store: &Arc>, ack: &Arc>, channel: &str, row_id: i64) { diff --git a/src/gateway/tests.rs b/src/gateway/tests.rs index 1475abd..674f3ee 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -208,7 +208,6 @@ fn setup_failure_ctx( .unwrap(); assert_eq!(inbound_id, 1); Ctx { - telegram_reply_anchor: std::sync::Arc::new(std::sync::Mutex::new(None)), cfg: test_config( &temp_path("setup-failure-state").to_string_lossy(), &temp_path("setup-failure-sessions").to_string_lossy(), @@ -399,7 +398,7 @@ async fn delivery_fails_when_channel_produces_no_chunks() { )) .unwrap(); - assert!(!reply_to(&gateway.ctx, "me@icloud.com", " \t\n ").await); + assert!(!reply_to(&gateway.ctx, "me@icloud.com", " \t\n ", None).await); assert!(gateway.ctx.sent_replies.lock().unwrap().is_empty()); let checkpoints = Arc::new(Mutex::new(Vec::new())); diff --git a/src/gateway/worker.rs b/src/gateway/worker.rs index 28bb16e..81041be 100644 --- a/src/gateway/worker.rs +++ b/src/gateway/worker.rs @@ -576,8 +576,13 @@ async fn review_changed_schedules(ctx: &Ctx, job: &Job) { }; audit_schedule_events(ctx, &mut ledger); for question in questions { - let delivered = - super::reply_to(ctx, &question.target, &question.render_correlated_text()).await; + let delivered = super::reply_to( + ctx, + &question.target, + &question.render_correlated_text(), + None, + ) + .await; if let Err(error) = ledger.mark_schedule_question_delivery( &question.id, if delivered { @@ -848,7 +853,6 @@ pub(super) async fn record_and_deliver( origin: OutboundOrigin, text: &str, ) -> Result { - *ctx.telegram_reply_anchor.lock().unwrap() = job.telegram_reply_anchor; let outbound = ctx.history.lock().unwrap().record_outbound( job.inbound_id, origin, @@ -875,7 +879,8 @@ pub(super) async fn record_and_deliver_once( if outbound.status == DeliveryStatus::Delivered { return Ok(true); } - let delivered = deliver_outbound_once(ctx, &job.target, &mut outbound).await?; + let delivered = + deliver_outbound_once(ctx, &job.target, job.telegram_reply_anchor, &mut outbound).await?; ctx.history.lock().unwrap().mark_delivery( outbound.id, if delivered { @@ -900,7 +905,9 @@ async fn deliver_stored( let mut attempt = 0; loop { attempt += 1; - let delivered = deliver_outbound_once(ctx, &job.target, &mut outbound).await?; + let delivered = + deliver_outbound_once(ctx, &job.target, job.telegram_reply_anchor, &mut outbound) + .await?; let status = if delivered { DeliveryStatus::Delivered } else { @@ -945,6 +952,7 @@ async fn deliver_stored( async fn deliver_outbound_once( ctx: &Ctx, target: &str, + reply_anchor: Option, outbound: &mut OutboundMessage, ) -> Result { let chunks = ctx @@ -963,7 +971,7 @@ async fn deliver_outbound_once( .enumerate() .skip(outbound.delivery_chunk_index) { - if let Err(error) = super::send_reply_chunk(ctx, target, &mut chunk.clone()).await { + if let Err(error) = super::send_reply_chunk(ctx, target, reply_anchor, chunk).await { error!( "outbound {} chunk {index} send error to {target}: {error}", outbound.id diff --git a/src/telegram.rs b/src/telegram.rs index b0fe522..3057bfc 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -280,7 +280,8 @@ impl Telegram { payload["text"] = json!(html); payload["parse_mode"] = json!("HTML"); if let Some(id) = reply_to_message_id { - payload["reply_parameters"] = json!({"message_id": id, "allow_sending_without_reply": true}); + payload["reply_parameters"] = + json!({"message_id": id, "allow_sending_without_reply": true}); } let transport_response = self .post_with_topic_fallback("sendMessage", payload) @@ -309,7 +310,8 @@ impl Telegram { let mut payload = target_payload(target); payload["text"] = json!(text); if let Some(id) = reply_to_message_id { - payload["reply_parameters"] = json!({"message_id": id, "allow_sending_without_reply": true}); + payload["reply_parameters"] = + json!({"message_id": id, "allow_sending_without_reply": true}); } let transport_response = self .post_with_topic_fallback("sendMessage", payload) @@ -847,6 +849,29 @@ mod tests { assert!(!messages[3].is_group); } + #[tokio::test] + async fn replies_carry_reply_parameters_only_when_anchored() { + let ok = json!({"ok": true, "result": {"message_id": 5}}); + let fake = Arc::new(FakeTransport::with_responses(vec![ok.clone(), ok])); + let telegram = + Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); + + telegram.send_plain("chat", "hi").await.unwrap(); + telegram + .send_plain_reply("chat", "hi", Some(41)) + .await + .unwrap(); + + let calls = fake.calls.lock().unwrap(); + assert_eq!(calls.len(), 2); + assert!(calls[0].1.get("reply_parameters").is_none()); + assert_eq!(calls[1].1["reply_parameters"]["message_id"], 41); + assert_eq!( + calls[1].1["reply_parameters"]["allow_sending_without_reply"], + true + ); + } + #[tokio::test] async fn poll_uses_next_update_offset_and_long_poll_timeout() { let fake = Arc::new(FakeTransport::with_responses(vec![json!({ From 6387a49b6c5ec4372734fb81210e8caa12e351eb Mon Sep 17 00:00:00 2001 From: Bernardo Kuri Date: Thu, 17 Sep 2026 18:21:46 -0600 Subject: [PATCH 05/10] feat(telegram): surface inbound reply-to context to the backend Telegram updates carry reply_to_message when the user replies to a specific message; the gateway dropped it at parse time, so the backend saw only the bare reply text with no linkage. Capture the replied-to message as a flattened one-line excerpt (200 char cap) on RawMessage.reply_context and prefix it to the backend prompt as a blockquote. The quote is context only: approval answer matching and the /stop check still see the raw user text, so replying "/stop" to an old message still stops the run. iMessage and Slack set reply_context: None. --- src/audit.rs | 1 + src/channel.rs | 7 +++++++ src/gateway/mod.rs | 13 +++++++++--- src/gateway/tests.rs | 4 ++++ src/slack.rs | 1 + src/telegram.rs | 50 ++++++++++++++++++++++++++++++++++++++++++++ 6 files changed, 73 insertions(+), 3 deletions(-) diff --git a/src/audit.rs b/src/audit.rs index 5c01aa8..34132a9 100644 --- a/src/audit.rs +++ b/src/audit.rs @@ -396,6 +396,7 @@ mod tests { is_from_me: false, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, } } diff --git a/src/channel.rs b/src/channel.rs index 80a0985..4facad5 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -58,6 +58,9 @@ pub struct RawMessage { /// Provider message id of the inbound message (Telegram `message_id`), /// used to anchor replies via `reply_parameters`. pub reply_to_message_id: Option, + /// Flattened excerpt of the message the user replied to (Telegram + /// `reply_to_message`), prefixed to the backend prompt. None elsewhere. + pub reply_context: Option, } impl RawMessage { @@ -428,6 +431,7 @@ impl ChannelContract for IMessageChannel { is_from_me: message.is_from_me, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, }) .collect()) @@ -895,6 +899,7 @@ mod tests { is_from_me, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -913,6 +918,7 @@ mod tests { is_from_me: false, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -1011,6 +1017,7 @@ mod tests { is_from_me: false, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, }; diff --git a/src/gateway/mod.rs b/src/gateway/mod.rs index 2484725..9e16662 100644 --- a/src/gateway/mod.rs +++ b/src/gateway/mod.rs @@ -894,21 +894,28 @@ impl Gateway { "[{thread}] new message accepted; routing to {}", backend.as_str() ); + // The reply quote is context for the backend only: approval + // answer matching and the /stop check below must see the raw + // user text. + let is_stop = + m.images.is_empty() && message_text.trim().eq_ignore_ascii_case("/stop"); let job = Job { row_id: m.row_id, inbound_id, thread, target, backend, - text: message_text, + text: match &m.reply_context { + Some(quoted) => format!("> {quoted}\n{message_text}"), + None => message_text, + }, reply_with_voice, voice_attachment: m.voice.clone(), image_attachments: m.images.clone(), approval_origin, telegram_reply_anchor: m.reply_to_message_id, }; - if job.image_attachments.is_empty() && job.text.trim().eq_ignore_ascii_case("/stop") - { + if is_stop { if !self.stop(job).await { return; } diff --git a/src/gateway/tests.rs b/src/gateway/tests.rs index 674f3ee..b2bb460 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -167,6 +167,7 @@ fn msg(chat: &str, handle: &str, from_me: bool, text: &str) -> RawMessage { is_group: false, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -4101,6 +4102,7 @@ fn message(row_id: i64, chat: &str, handle: &str, is_from_me: bool, text: &str) is_group: false, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -4125,6 +4127,7 @@ fn telegram_message( is_group, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -4180,6 +4183,7 @@ fn slack_image_message( is_group: false, is_supported: true, reply_to_message_id: None, + reply_context: None, thread_id: None, } } diff --git a/src/slack.rs b/src/slack.rs index 7a2296c..778945c 100644 --- a/src/slack.rs +++ b/src/slack.rs @@ -664,6 +664,7 @@ impl Inbox { is_from_me: row.get(8)?, is_supported: row.get(9)?, reply_to_message_id: None, + reply_context: None, thread_id: None, }) })? diff --git a/src/telegram.rs b/src/telegram.rs index 3057bfc..ff6e65d 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -544,6 +544,7 @@ impl Update { is_supported: false, thread_id: None, reply_to_message_id: None, + reply_context: None, }; }; let images = message @@ -600,10 +601,26 @@ impl Update { is_supported: true, thread_id: message.message_thread_id, reply_to_message_id: message.message_id, + reply_context: message + .reply_to_message + .as_deref() + .and_then(|replied| replied.text.as_deref().or(replied.caption.as_deref())) + .map(reply_excerpt), } } } +/// Flattens a replied-to message to one line, capped so a long quote cannot +/// dominate the prompt. +fn reply_excerpt(text: &str) -> String { + let flat = text.split_whitespace().collect::>().join(" "); + let mut excerpt: String = flat.chars().take(200).collect(); + if flat.chars().count() > 200 { + excerpt.push('…'); + } + excerpt +} + #[derive(Deserialize)] struct TelegramMessage { #[serde(default)] @@ -623,6 +640,8 @@ struct TelegramMessage { voice: Option, #[serde(default)] message_thread_id: Option, + #[serde(default)] + reply_to_message: Option>, } #[derive(Deserialize)] @@ -849,6 +868,35 @@ mod tests { assert!(!messages[3].is_group); } + #[test] + fn parses_reply_to_message_into_flattened_excerpt() { + let long = "word ".repeat(80).trim_end().to_string(); + let update: Update = serde_json::from_value(json!({ + "update_id": 9, + "message": { + "message_id": 20, + "from": {"id": 7}, + "chat": {"id": 7, "type": "private"}, + "text": "now do the other thing", + "reply_to_message": { + "message_id": 19, + "from": {"id": 99}, + "chat": {"id": 7, "type": "private"}, + "text": format!("first line\nsecond line {long}") + } + } + })) + .unwrap(); + + let message = update.into_raw(); + + let quoted = message.reply_context.unwrap(); + assert!(!quoted.contains('\n')); + assert!(quoted.starts_with("first line second line word")); + assert!(quoted.ends_with('…')); + assert_eq!(quoted.chars().count(), 201); + } + #[tokio::test] async fn replies_carry_reply_parameters_only_when_anchored() { let ok = json!({"ok": true, "result": {"message_id": 5}}); @@ -921,6 +969,7 @@ mod tests { voice: None, message_thread_id: None, message_id: Some(1), + reply_to_message: None, }), } .into_raw(); @@ -953,6 +1002,7 @@ mod tests { }), message_thread_id: None, message_id: Some(1), + reply_to_message: None, }), } .into_raw(); From add3b7246608f5e570a186714efbcfa371f65782 Mon Sep 17 00:00:00 2001 From: Bernardo Kuri Date: Thu, 17 Sep 2026 18:55:44 -0600 Subject: [PATCH 06/10] feat(telegram): prefer the user's selected quote for reply context When the user selects a specific portion of a message and replies, Telegram sends the selection in message.quote. Use it for the backend excerpt when present; fall back to the full replied-to message text. --- src/telegram.rs | 52 +++++++++++++++++++++++++++++++++++++++++++++---- 1 file changed, 48 insertions(+), 4 deletions(-) diff --git a/src/telegram.rs b/src/telegram.rs index ff6e65d..38e794c 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -602,10 +602,16 @@ impl Update { thread_id: message.message_thread_id, reply_to_message_id: message.message_id, reply_context: message - .reply_to_message - .as_deref() - .and_then(|replied| replied.text.as_deref().or(replied.caption.as_deref())) - .map(reply_excerpt), + .quote + .map(|quote| quote.text) + .or_else(|| { + message + .reply_to_message + .as_deref() + .and_then(|replied| replied.text.as_deref().or(replied.caption.as_deref())) + .map(str::to_string) + }) + .map(|text| reply_excerpt(&text)), } } } @@ -641,9 +647,17 @@ struct TelegramMessage { #[serde(default)] message_thread_id: Option, #[serde(default)] + quote: Option, + #[serde(default)] reply_to_message: Option>, } +#[derive(Deserialize)] +struct TextQuote { + #[serde(default)] + text: String, +} + #[derive(Deserialize)] struct TelegramVoice { file_id: String, @@ -868,6 +882,34 @@ mod tests { assert!(!messages[3].is_group); } + #[test] + fn selected_quote_takes_priority_over_replied_to_text() { + let update: Update = serde_json::from_value(json!({ + "update_id": 10, + "message": { + "message_id": 21, + "from": {"id": 7}, + "chat": {"id": 7, "type": "private"}, + "text": "what about this part", + "quote": {"position": 12, "text": "the frobnicator keeps timing out"}, + "reply_to_message": { + "message_id": 20, + "from": {"id": 99}, + "chat": {"id": 7, "type": "private"}, + "text": "long full message the user did not select" + } + } + })) + .unwrap(); + + let message = update.into_raw(); + + assert_eq!( + message.reply_context.as_deref(), + Some("the frobnicator keeps timing out") + ); + } + #[test] fn parses_reply_to_message_into_flattened_excerpt() { let long = "word ".repeat(80).trim_end().to_string(); @@ -969,6 +1011,7 @@ mod tests { voice: None, message_thread_id: None, message_id: Some(1), + quote: None, reply_to_message: None, }), } @@ -1002,6 +1045,7 @@ mod tests { }), message_thread_id: None, message_id: Some(1), + quote: None, reply_to_message: None, }), } From 8dbd94271b5026eb44b41fee82eff0cef80c7681 Mon Sep 17 00:00:00 2001 From: Bernardo Kuri Date: Thu, 17 Sep 2026 19:08:15 -0600 Subject: [PATCH 07/10] feat(telegram): backend can quote the exact portion it is answering A backend reply may start with a '> quoted portion' line; the gateway lifts it into reply_parameters.quote (1024-char cap, API limit), strips it from the displayed text, and it renders in the Telegram reply header when the message is anchored. Quote-only replies are left untouched so they cannot become empty messages. System and scheduled sends never quote. --- src/channel.rs | 26 +++++++++++++++++---- src/gateway/mod.rs | 30 ++++++++++++++++++++++-- src/gateway/tests.rs | 20 ++++++++++++++++ src/gateway/worker.rs | 36 ++++++++++++++++++++--------- src/telegram.rs | 54 +++++++++++++++++++++++++++++++++++-------- 5 files changed, 140 insertions(+), 26 deletions(-) diff --git a/src/channel.rs b/src/channel.rs index 4facad5..4b29611 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -78,6 +78,9 @@ pub struct OutboundChunk { pub rich_markdown: bool, /// Telegram message id to reply to (anchoring), when known. pub reply_to_message_id: Option, + /// Exact portion of the trigger message to quote in the reply header + /// (Telegram `reply_parameters.quote`), chosen by the backend. + pub quote: Option, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -526,6 +529,7 @@ impl ChannelContract for IMessageChannel { } else { vec![OutboundChunk { reply_to_message_id: None, + quote: None, text: format!("{text}{marker}"), rich_markdown: false, }] @@ -642,6 +646,7 @@ impl ChannelContract for Telegram { .into_iter() .map(|text| OutboundChunk { reply_to_message_id: None, + quote: None, text, rich_markdown: true, }) @@ -650,11 +655,21 @@ impl ChannelContract for Telegram { async fn send_chunk(&self, target: &str, chunk: &OutboundChunk) -> Result<()> { if chunk.rich_markdown { - self.send_rich_reply(target, &chunk.text, chunk.reply_to_message_id) - .await + self.send_rich_reply( + target, + &chunk.text, + chunk.reply_to_message_id, + chunk.quote.as_deref(), + ) + .await } else { - self.send_plain_reply(target, &chunk.text, chunk.reply_to_message_id) - .await + self.send_plain_reply( + target, + &chunk.text, + chunk.reply_to_message_id, + chunk.quote.as_deref(), + ) + .await } } @@ -751,6 +766,7 @@ impl ChannelContract for Slack { .into_iter() .map(|text| OutboundChunk { reply_to_message_id: None, + quote: None, text, rich_markdown: true, }) @@ -959,6 +975,7 @@ mod tests { channel.outbound_chunks("reply", REPLY_MARKER), [OutboundChunk { reply_to_message_id: None, + quote: None, text: format!("reply{REPLY_MARKER}"), rich_markdown: false, }] @@ -995,6 +1012,7 @@ mod tests { channel.outbound_chunks("reply", REPLY_MARKER), [OutboundChunk { reply_to_message_id: None, + quote: None, text: "reply".to_string(), rich_markdown: true, }] diff --git a/src/gateway/mod.rs b/src/gateway/mod.rs index 9e16662..efc6cde 100644 --- a/src/gateway/mod.rs +++ b/src/gateway/mod.rs @@ -1237,6 +1237,30 @@ fn audit_schedule_events(ctx: &Ctx, ledger: &mut jobs::Ledger) { } } +/// Lifts a leading `> quoted portion` line from a backend reply into the +/// Telegram reply quote and strips it from the displayed text. A reply that +/// is only a quote line is left untouched so it cannot become empty. +fn lift_outbound_quote(text: &str) -> (Option, String) { + let trimmed = text.trim_start(); + let Some(first_line) = trimmed.lines().next() else { + return (None, text.to_string()); + }; + let Some(quoted) = first_line + .strip_prefix("> ") + .map(str::trim) + .filter(|q| !q.is_empty()) + else { + return (None, text.to_string()); + }; + let rest = trimmed[first_line.len()..] + .trim_start_matches(['\r', '\n']) + .to_string(); + if rest.is_empty() { + return (None, text.to_string()); + } + (Some(quoted.to_string()), rest) +} + async fn reply_to(ctx: &Ctx, target: &str, text: &str, reply_anchor: Option) -> bool { let chunks = ctx.channel.outbound_chunks(text, &ctx.reply_marker); if chunks.is_empty() { @@ -1244,7 +1268,7 @@ async fn reply_to(ctx: &Ctx, target: &str, text: &str, reply_anchor: Option return false; } for chunk in chunks { - if let Err(error) = send_reply_chunk(ctx, target, reply_anchor, &chunk).await { + if let Err(error) = send_reply_chunk(ctx, target, reply_anchor, None, &chunk).await { error!("send error to {target}: {error}"); return false; } @@ -1257,6 +1281,7 @@ async fn send_reply_chunk( ctx: &Ctx, target: &str, reply_anchor: Option, + quote: Option<&str>, chunk: &crate::channel::OutboundChunk, ) -> Result<()> { #[cfg(test)] @@ -1294,6 +1319,7 @@ async fn send_reply_chunk( { let mut chunk = chunk.clone(); chunk.reply_to_message_id = reply_anchor; + chunk.quote = quote.map(str::to_string); let timeout = ctx.channel.delivery_semantics().send_timeout; if timeout.is_zero() { ctx.channel.send_chunk(target, &chunk).await @@ -1342,7 +1368,7 @@ async fn send_scheduled_chunk( target: &str, chunk: &crate::channel::OutboundChunk, ) -> Result<()> { - send_reply_chunk(ctx, target, None, chunk).await + send_reply_chunk(ctx, target, None, None, chunk).await } fn complete_row(store: &Arc>, ack: &Arc>, channel: &str, row_id: i64) { diff --git a/src/gateway/tests.rs b/src/gateway/tests.rs index b2bb460..d265999 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -4187,3 +4187,23 @@ fn slack_image_message( thread_id: None, } } + +#[test] +fn lifts_leading_quote_line_from_backend_replies() { + let (quote, rest) = + super::lift_outbound_quote("> the frobnicator keeps timing out\nTry the cache path."); + assert_eq!(quote.as_deref(), Some("the frobnicator keeps timing out")); + assert_eq!(rest, "Try the cache path."); + + let (quote, rest) = super::lift_outbound_quote("No quote here.\n> not lifted"); + assert!(quote.is_none()); + assert_eq!(rest, "No quote here.\n> not lifted"); + + let (quote, rest) = super::lift_outbound_quote("> quote only"); + assert!(quote.is_none()); + assert_eq!(rest, "> quote only"); + + let (quote, rest) = super::lift_outbound_quote("> \nafter empty quote"); + assert!(quote.is_none()); + assert_eq!(rest, "> \nafter empty quote"); +} diff --git a/src/gateway/worker.rs b/src/gateway/worker.rs index 81041be..eb1aabe 100644 --- a/src/gateway/worker.rs +++ b/src/gateway/worker.rs @@ -68,7 +68,7 @@ where report_delivery( ctx, &job, - deliver_stored(ctx, &job, &outbound).await, + deliver_stored(ctx, &job, &outbound, None).await, &outbound.content, "recovered_outbound", "recover outbound", @@ -443,7 +443,7 @@ where ); return; } - let delivery = deliver_stored(ctx, &job, &outbound).await; + let delivery = deliver_stored(ctx, &job, &outbound, None).await; if delivery.is_ok() { info!("[{}] reply sent via {}", job.thread, ctx.channel.id()); } @@ -718,7 +718,7 @@ async fn finish_run_with_gateway_reply(ctx: &Ctx, job: &Job, reply: &str, labels report_delivery( ctx, job, - deliver_stored(ctx, job, &outbound).await, + deliver_stored(ctx, job, &outbound, None).await, reply, labels.completion, labels.deliver, @@ -853,13 +853,14 @@ pub(super) async fn record_and_deliver( origin: OutboundOrigin, text: &str, ) -> Result { + let (quote, text) = super::lift_outbound_quote(text); let outbound = ctx.history.lock().unwrap().record_outbound( job.inbound_id, origin, Some(job.backend.as_str()), - text, + &text, )?; - deliver_stored(ctx, job, &outbound).await + deliver_stored(ctx, job, &outbound, quote.as_deref()).await } /// Records an outbound reply and tries delivery once. Control-path replies use @@ -879,8 +880,14 @@ pub(super) async fn record_and_deliver_once( if outbound.status == DeliveryStatus::Delivered { return Ok(true); } - let delivered = - deliver_outbound_once(ctx, &job.target, job.telegram_reply_anchor, &mut outbound).await?; + let delivered = deliver_outbound_once( + ctx, + &job.target, + job.telegram_reply_anchor, + None, + &mut outbound, + ) + .await?; ctx.history.lock().unwrap().mark_delivery( outbound.id, if delivered { @@ -896,6 +903,7 @@ async fn deliver_stored( ctx: &Ctx, job: &Job, outbound: &OutboundMessage, + quote: Option<&str>, ) -> Result { if outbound.status == DeliveryStatus::Delivered { return Ok(DeliveryOutcome::AlreadyDelivered); @@ -905,9 +913,14 @@ async fn deliver_stored( let mut attempt = 0; loop { attempt += 1; - let delivered = - deliver_outbound_once(ctx, &job.target, job.telegram_reply_anchor, &mut outbound) - .await?; + let delivered = deliver_outbound_once( + ctx, + &job.target, + job.telegram_reply_anchor, + quote, + &mut outbound, + ) + .await?; let status = if delivered { DeliveryStatus::Delivered } else { @@ -953,6 +966,7 @@ async fn deliver_outbound_once( ctx: &Ctx, target: &str, reply_anchor: Option, + quote: Option<&str>, outbound: &mut OutboundMessage, ) -> Result { let chunks = ctx @@ -971,7 +985,7 @@ async fn deliver_outbound_once( .enumerate() .skip(outbound.delivery_chunk_index) { - if let Err(error) = super::send_reply_chunk(ctx, target, reply_anchor, chunk).await { + if let Err(error) = super::send_reply_chunk(ctx, target, reply_anchor, quote, chunk).await { error!( "outbound {} chunk {index} send error to {target}: {error}", outbound.id diff --git a/src/telegram.rs b/src/telegram.rs index 38e794c..b5e3a0f 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -263,7 +263,7 @@ impl Telegram { } pub async fn send_rich(&self, target: &str, text: &str) -> Result<()> { - self.send_rich_reply(target, text, None).await + self.send_rich_reply(target, text, None, None).await } pub async fn send_rich_reply( @@ -271,6 +271,7 @@ impl Telegram { target: &str, text: &str, reply_to_message_id: Option, + quote: Option<&str>, ) -> Result<()> { if text.encode_utf16().count() > TEXT_LIMIT { bail!("Telegram rich message exceeds the {TEXT_LIMIT} character chunk limit"); @@ -279,9 +280,8 @@ impl Telegram { let mut payload = target_payload(target); payload["text"] = json!(html); payload["parse_mode"] = json!("HTML"); - if let Some(id) = reply_to_message_id { - payload["reply_parameters"] = - json!({"message_id": id, "allow_sending_without_reply": true}); + if let Some(params) = reply_parameters_payload(reply_to_message_id, quote) { + payload["reply_parameters"] = params; } let transport_response = self .post_with_topic_fallback("sendMessage", payload) @@ -298,7 +298,7 @@ impl Telegram { } pub async fn send_plain(&self, target: &str, text: &str) -> Result<()> { - self.send_plain_reply(target, text, None).await + self.send_plain_reply(target, text, None, None).await } pub async fn send_plain_reply( @@ -306,12 +306,12 @@ impl Telegram { target: &str, text: &str, reply_to_message_id: Option, + quote: Option<&str>, ) -> Result<()> { let mut payload = target_payload(target); payload["text"] = json!(text); - if let Some(id) = reply_to_message_id { - payload["reply_parameters"] = - json!({"message_id": id, "allow_sending_without_reply": true}); + if let Some(params) = reply_parameters_payload(reply_to_message_id, quote) { + payload["reply_parameters"] = params; } let transport_response = self .post_with_topic_fallback("sendMessage", payload) @@ -527,6 +527,23 @@ struct Update { message: Option, } +/// Builds `reply_parameters` for anchored sends. The quote (a portion of the +/// replied-to message chosen by the backend) only renders when anchored, and +/// Telegram caps it at 1024 characters. +fn reply_parameters_payload( + reply_to_message_id: Option, + quote: Option<&str>, +) -> Option { + let mut params = json!({ + "message_id": reply_to_message_id?, + "allow_sending_without_reply": true + }); + if let Some(quoted) = quote { + params["quote"] = json!(quoted.chars().take(1024).collect::()); + } + Some(params) +} + impl Update { fn into_raw(self) -> RawMessage { let Some(message) = self.message else { @@ -948,7 +965,7 @@ mod tests { telegram.send_plain("chat", "hi").await.unwrap(); telegram - .send_plain_reply("chat", "hi", Some(41)) + .send_plain_reply("chat", "hi", Some(41), Some("the frobnicator")) .await .unwrap(); @@ -960,6 +977,25 @@ mod tests { calls[1].1["reply_parameters"]["allow_sending_without_reply"], true ); + assert_eq!(calls[1].1["reply_parameters"]["quote"], "the frobnicator"); + } + + #[tokio::test] + async fn quote_is_dropped_without_an_anchor() { + let fake = Arc::new(FakeTransport::with_responses(vec![json!({ + "ok": true, + "result": {"message_id": 5} + })])); + let telegram = + Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); + + telegram + .send_plain_reply("chat", "hi", None, Some("orphan quote")) + .await + .unwrap(); + + let calls = fake.calls.lock().unwrap(); + assert!(calls[0].1.get("reply_parameters").is_none()); } #[tokio::test] From 13831f47f67920a046d339a83de9401de4083112 Mon Sep 17 00:00:00 2001 From: Bernardo Kuri Date: Thu, 17 Sep 2026 19:40:17 -0600 Subject: [PATCH 08/10] docs(skill): document the reply-quote convention in the managed push skill The '> quoted portion' convention ships in the managed skill so every assistant learns it when push init refreshes the skill (managed version bumped to 4 to trigger the refresh on upgrade). --- assistant/skills/push/SKILL.md | 6 ++++++ src/assistant.rs | 2 +- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/assistant/skills/push/SKILL.md b/assistant/skills/push/SKILL.md index 98d60b2..9eb903b 100644 --- a/assistant/skills/push/SKILL.md +++ b/assistant/skills/push/SKILL.md @@ -80,3 +80,9 @@ in diagnostics or replies. For an ordinary conversation, return the final reply normally. Push sends it back through the originating channel. Do not invoke the Push CLI merely to send a chat reply. + +When the reply answers a specific part of the user's message, start the reply +with a single line `> ` before the answer. On channels that support +reply quotes, Push moves that line into the reply header quote and removes it +from the displayed text. Quote only exact words from the user's message, and +send no quote line when the reply addresses the whole message. diff --git a/src/assistant.rs b/src/assistant.rs index 0f88700..7f7902a 100644 --- a/src/assistant.rs +++ b/src/assistant.rs @@ -61,7 +61,7 @@ Good examples include preferences, active projects, people, recurring processes, "#; const PUSH_SKILL: &str = include_str!("../assistant/skills/push/SKILL.md"); -const PUSH_SKILL_VERSION: u32 = 3; +const PUSH_SKILL_VERSION: u32 = 4; const PUSH_SKILL_LINK: &str = "../../skills/push"; const PUSH_SKILL_MANIFEST: &str = ".push-managed.json"; const PUSH_SKILL_PROVIDERS: [&str; 2] = [".agents", ".claude"]; From a3cb8da847e6b3e63c83cfe489b95375b2bbeb84 Mon Sep 17 00:00:00 2001 From: Bernardo Kuri Date: Thu, 17 Sep 2026 22:31:04 -0600 Subject: [PATCH 09/10] fix(gateway): lift quotes on the backend reply path; keep anchor across rich fallback Greptile review follow-ups: - The backend run completion path recorded out.reply and delivered with no lifted quote, so the '> quoted portion' convention never fired for the one reply type that carries it. Lift before recording and thread the quote through deliver_stored. - The rich-to-plain fallback in send_rich_reply dropped the reply anchor and quote; pass them through to send_plain_reply. - Drop the now-unused send_rich/send_plain wrappers. - FakeRunner gains an optional reply override; new gateway test asserts a leading quote line is stripped from the delivered reply. --- src/agent.rs | 12 ++++++--- src/gateway/tests.rs | 61 +++++++++++++++++++++++++++++++++++++++++++ src/gateway/worker.rs | 7 ++--- src/telegram.rs | 61 +++++++++++++++++++++++++++++-------------- 4 files changed, 115 insertions(+), 26 deletions(-) diff --git a/src/agent.rs b/src/agent.rs index 529fff5..7d3db5d 100644 --- a/src/agent.rs +++ b/src/agent.rs @@ -177,6 +177,9 @@ pub struct FakeRunner { pub wait_for_release: Option>, pub failure: Option, pub resume_missing_once: Option>, + /// Overrides the canned reply so tests can exercise backend-reply + /// formatting conventions (for example a leading `> quoted portion`). + pub reply: Option, } #[cfg(test)] @@ -221,11 +224,14 @@ impl FakeRunner { return Err(RunError::Failed(message.clone())); } let current_message = crate::prompt::current_message(req.prompt); - Ok(RunOutput { - reply: format!( + let reply = self.reply.clone().unwrap_or_else(|| { + format!( "fake reply: {}", current_message.as_deref().unwrap_or(req.prompt) - ), + ) + }); + Ok(RunOutput { + reply, session_id: req.is_new.then(|| self.session_id.clone()), }) } diff --git a/src/gateway/tests.rs b/src/gateway/tests.rs index d265999..064fd04 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -1077,6 +1077,7 @@ async fn missing_backend_session_rotates_and_rehydrates_once() { wait_for_release: None, failure: None, resume_missing_once: Some(missing), + reply: None, }), ); gateway.ctx.runners = Arc::new(runners); @@ -1161,6 +1162,7 @@ async fn backend_switch_and_clear_start_fresh_sessions_with_history() { wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), ); runners.insert( @@ -1173,6 +1175,7 @@ async fn backend_switch_and_clear_start_fresh_sessions_with_history() { wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), ); gateway.ctx.runners = Arc::new(runners); @@ -1887,6 +1890,7 @@ async fn slack_images_reach_every_agent_backend_and_are_removed_after_each_turn( wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), )])); @@ -1946,6 +1950,7 @@ async fn imessage_images_reach_every_agent_backend_and_are_removed_after_each_tu wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), )])); let mut inbound = message(1, "+15551234567", "+15551234567", false, ""); @@ -2476,6 +2481,7 @@ async fn telegram_image_only_and_captioned_messages_reach_pi() { wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), )])); @@ -2636,6 +2642,7 @@ async fn telegram_image_reaches_claude_and_is_removed_after_the_turn() { wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), )])); run_messages( @@ -2876,6 +2883,7 @@ async fn stop_interrupts_active_run_and_preserves_queued_messages() { wait_for_release: Some(release.clone()), failure: None, resume_missing_once: None, + reply: None, }), ); gateway.ctx.runners = Arc::new(runners); @@ -2985,6 +2993,7 @@ async fn stop_interrupts_a_worker_queued_in_the_same_poll_batch() { wait_for_release: Some(release.clone()), failure: None, resume_missing_once: None, + reply: None, }), ); gateway.ctx.runners = Arc::new(runners); @@ -3253,6 +3262,7 @@ async fn retried_stop_acknowledgement_does_not_cancel_the_next_request() { wait_for_release: Some(release.clone()), failure: None, resume_missing_once: None, + reply: None, }), ); gateway.ctx.runners = Arc::new(runners); @@ -3318,6 +3328,7 @@ async fn retried_stop_acknowledgement_does_not_cancel_the_next_request() { wait_for_release: Some(restart_release.clone()), failure: None, resume_missing_once: None, + reply: None, }), ); restarted.ctx.runners = Arc::new(restart_runners); @@ -4040,6 +4051,7 @@ fn fake_runners_with_hook( wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), ); runners @@ -4207,3 +4219,52 @@ fn lifts_leading_quote_line_from_backend_replies() { assert!(quote.is_none()); assert_eq!(rest, "> \nafter empty quote"); } + +#[tokio::test] +async fn backend_reply_leading_quote_line_is_lifted_into_the_reply_quote() { + let state_path = temp_state_path(); + let sessions_dir = temp_path("backend-quote-sessions"); + let assistant_dir = temp_path("backend-quote-assistant"); + std::fs::create_dir_all(&assistant_dir).unwrap(); + let calls = Arc::new(Mutex::new(Vec::new())); + let config = test_config( + &state_path, + sessions_dir.to_str().unwrap(), + assistant_dir.to_str().unwrap(), + ); + let mut gateway = Gateway::new(config).unwrap(); + { + let mut runners = HashMap::new(); + runners.insert( + AgentBackend::Codex, + Runner::Fake(FakeRunner { + backend: AgentBackend::Codex, + session_id: "fake-session".to_string(), + calls: calls.clone(), + before_return: None, + wait_for_release: None, + failure: None, + resume_missing_once: None, + reply: Some("> the frobnicator keeps timing out\nTry the cache path.".to_string()), + }), + ); + gateway.ctx.runners = Arc::new(runners); + } + run_messages( + &mut gateway, + vec![message(1, "me@icloud.com", "", true, "hello")], + ) + .await; + + assert_eq!( + gateway.ctx.sent_replies.lock().unwrap().as_slice(), + [( + "me@icloud.com".to_string(), + "Try the cache path.\n\n-- sent by push".to_string() + )] + ); + + let _ = std::fs::remove_file(&state_path); + let _ = std::fs::remove_file(format!("{state_path}.db")); + let _ = std::fs::remove_file(format!("{state_path}.audit.jsonl")); +} diff --git a/src/gateway/worker.rs b/src/gateway/worker.rs index eb1aabe..ac7a0d5 100644 --- a/src/gateway/worker.rs +++ b/src/gateway/worker.rs @@ -412,11 +412,12 @@ where ctx.audit .backend_completed(job.row_id, &job.thread, job.backend, &out.reply), ); + let (reply_quote, reply_text) = super::lift_outbound_quote(&out.reply); let outbound = match ctx.history.lock().unwrap().record_outbound( job.inbound_id, OutboundOrigin::Backend, Some(job.backend.as_str()), - &out.reply, + &reply_text, ) { Ok(outbound) => outbound, Err(error) => { @@ -443,7 +444,7 @@ where ); return; } - let delivery = deliver_stored(ctx, &job, &outbound, None).await; + let delivery = deliver_stored(ctx, &job, &outbound, reply_quote.as_deref()).await; if delivery.is_ok() { info!("[{}] reply sent via {}", job.thread, ctx.channel.id()); } @@ -451,7 +452,7 @@ where ctx, &job, delivery, - &out.reply, + &reply_text, "completed", "deliver backend reply", ); diff --git a/src/telegram.rs b/src/telegram.rs index b5e3a0f..d104f9a 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -262,10 +262,6 @@ impl Telegram { self.allow_user_ids.contains(&chat_id) || self.allow_chat_ids.contains(&chat_id) } - pub async fn send_rich(&self, target: &str, text: &str) -> Result<()> { - self.send_rich_reply(target, text, None, None).await - } - pub async fn send_rich_reply( &self, target: &str, @@ -291,16 +287,14 @@ impl Telegram { if !response.ok { // Rendered HTML Telegram rejects (for example a parse error) // still reaches the user as plain text within the same durable - // delivery chunk. - self.send_plain(target, text).await?; + // delivery chunk. Keep the anchoring and quote across the + // fallback. + self.send_plain_reply(target, text, reply_to_message_id, quote) + .await?; } Ok(()) } - pub async fn send_plain(&self, target: &str, text: &str) -> Result<()> { - self.send_plain_reply(target, text, None, None).await - } - pub async fn send_plain_reply( &self, target: &str, @@ -963,7 +957,10 @@ mod tests { let telegram = Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); - telegram.send_plain("chat", "hi").await.unwrap(); + telegram + .send_plain_reply("chat", "hi", None, None) + .await + .unwrap(); telegram .send_plain_reply("chat", "hi", Some(41), Some("the frobnicator")) .await @@ -1272,7 +1269,10 @@ mod tests { let telegram = Telegram::with_transport("do-not-log".to_string(), vec![7], vec![], fake.clone()); - telegram.send_rich("7", "**reply**").await.unwrap(); + telegram + .send_rich_reply("7", "**reply**", None, None) + .await + .unwrap(); let calls = fake.calls.lock().unwrap(); assert_eq!( @@ -1299,7 +1299,10 @@ mod tests { Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); let text = "**reply**"; - telegram.send_rich("7", text).await.unwrap(); + telegram + .send_rich_reply("7", text, None, None) + .await + .unwrap(); let calls = fake.calls.lock().unwrap(); // The HTML send and its plain fallback share one gateway-owned chunk. @@ -1317,7 +1320,7 @@ mod tests { Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); let error = telegram - .send_rich("7", &"x".repeat(TEXT_LIMIT + 1)) + .send_rich_reply("7", &"x".repeat(TEXT_LIMIT + 1), None, None) .await .unwrap_err(); @@ -1335,8 +1338,14 @@ mod tests { let telegram = Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); - telegram.send_plain("7:99", "reply").await.unwrap(); - telegram.send_rich("7:99", "reply").await.unwrap(); + telegram + .send_plain_reply("7:99", "reply", None, None) + .await + .unwrap(); + telegram + .send_rich_reply("7:99", "reply", None, None) + .await + .unwrap(); telegram.send_typing("7:99").await.unwrap(); let calls = fake.calls.lock().unwrap(); @@ -1375,7 +1384,10 @@ mod tests { let telegram = Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); - telegram.send_plain("7:99", "reply").await.unwrap(); + telegram + .send_plain_reply("7:99", "reply", None, None) + .await + .unwrap(); let calls = fake.calls.lock().unwrap(); assert_eq!( @@ -1402,7 +1414,10 @@ mod tests { let telegram = Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); - let error = telegram.send_plain("7:99", "reply").await.unwrap_err(); + let error = telegram + .send_plain_reply("7:99", "reply", None, None) + .await + .unwrap_err(); assert!(error.to_string().contains("HTTP 400")); assert_eq!(fake.calls.lock().unwrap().len(), 1); @@ -1414,7 +1429,10 @@ mod tests { let telegram = Telegram::with_transport("secret".to_string(), vec![7], vec![], fake.clone()); - telegram.send_rich("7:99", "reply").await.unwrap(); + telegram + .send_rich_reply("7:99", "reply", None, None) + .await + .unwrap(); let calls = fake.calls.lock().unwrap(); assert_eq!(calls[0].0, "sendMessage"); @@ -1471,7 +1489,10 @@ mod tests { let text = format!("{}é", "a".repeat(TEXT_LIMIT)); for chunk in split_text(&text) { - telegram.send_plain("7", &chunk).await.unwrap(); + telegram + .send_plain_reply("7", &chunk, None, None) + .await + .unwrap(); } let calls = fake.calls.lock().unwrap(); From d559442ea6e63401f49f9a5995bc880138eaa87e Mon Sep 17 00:00:00 2001 From: Bernardo Kuri Date: Thu, 17 Sep 2026 22:42:17 -0600 Subject: [PATCH 10/10] fix(gateway): limit reply-quote lifting to Telegram Non-Telegram channels never render the extracted quote, so lifting a leading '> quoted' line from a backend reply there silently dropped it from the user's message. Gate the lift on the channel; Slack and iMessage keep the blockquote line as ordinary formatting. --- src/channel.rs | 6 ++++ src/gateway/tests.rs | 64 +++++++++++++++++++++++++++++++++++++++++-- src/gateway/worker.rs | 6 +++- 3 files changed, 72 insertions(+), 4 deletions(-) diff --git a/src/channel.rs b/src/channel.rs index 4b29611..4ffebd1 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -225,6 +225,12 @@ impl Channel { } } + /// Whether the channel renders `reply_parameters`-style reply quotes. + /// Quote lifting on other channels would silently drop the quoted line. + pub fn supports_reply_quotes(&self) -> bool { + matches!(self, Self::Telegram(_)) + } + pub async fn poll(&self, since: i64) -> Result> { match self { Self::IMessage(channel) => ChannelContract::poll(channel, since).await, diff --git a/src/gateway/tests.rs b/src/gateway/tests.rs index 064fd04..1a1888e 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -4220,8 +4220,8 @@ fn lifts_leading_quote_line_from_backend_replies() { assert_eq!(rest, "> \nafter empty quote"); } -#[tokio::test] -async fn backend_reply_leading_quote_line_is_lifted_into_the_reply_quote() { +#[tokio::test(flavor = "current_thread")] +async fn imessage_reply_keeps_leading_blockquote_line() { let state_path = temp_state_path(); let sessions_dir = temp_path("backend-quote-sessions"); let assistant_dir = temp_path("backend-quote-assistant"); @@ -4256,15 +4256,73 @@ async fn backend_reply_leading_quote_line_is_lifted_into_the_reply_quote() { ) .await; + // Non-Telegram channels do not render reply quotes, so the blockquote + // line must stay in the message text. assert_eq!( gateway.ctx.sent_replies.lock().unwrap().as_slice(), [( "me@icloud.com".to_string(), - "Try the cache path.\n\n-- sent by push".to_string() + "> the frobnicator keeps timing out\nTry the cache path.\n\n-- sent by push" + .to_string() )] ); let _ = std::fs::remove_file(&state_path); let _ = std::fs::remove_file(format!("{state_path}.db")); let _ = std::fs::remove_file(format!("{state_path}.audit.jsonl")); + let _ = std::fs::remove_dir_all(sessions_dir); + let _ = std::fs::remove_dir_all(assistant_dir); +} + +#[tokio::test(flavor = "current_thread")] +async fn telegram_backend_reply_leading_quote_line_is_lifted_into_the_reply_quote() { + let state_path = temp_state_path(); + let sessions_dir = temp_path("backend-quote-telegram-sessions"); + let assistant_dir = temp_path("backend-quote-telegram-assistant"); + std::fs::create_dir_all(&assistant_dir).unwrap(); + let calls = Arc::new(Mutex::new(Vec::new())); + let mut cfg = test_config( + &state_path, + sessions_dir.to_str().unwrap(), + assistant_dir.to_str().unwrap(), + ); + cfg.channel = "telegram".to_string(); + cfg.self_handles.clear(); + cfg.allow_from.clear(); + cfg.telegram_bot_token = Some("secret".to_string()); + cfg.telegram_allow_user_ids = vec![7]; + let mut gateway = Gateway::new(cfg).unwrap(); + { + let mut runners = HashMap::new(); + runners.insert( + AgentBackend::Codex, + Runner::Fake(FakeRunner { + backend: AgentBackend::Codex, + session_id: "fake-session".to_string(), + calls: calls.clone(), + before_return: None, + wait_for_release: None, + failure: None, + resume_missing_once: None, + reply: Some("> the frobnicator keeps timing out\nTry the cache path.".to_string()), + }), + ); + gateway.ctx.runners = Arc::new(runners); + } + gateway + .tick_fake(vec![telegram_message(1, 7, 7, false, "hello")]) + .await; + gateway.queues.clear(); + gateway.drain_workers().await; + + assert_eq!( + gateway.ctx.sent_replies.lock().unwrap().as_slice(), + [("7".to_string(), "Try the cache path.".to_string())] + ); + + let _ = std::fs::remove_file(&state_path); + let _ = std::fs::remove_file(format!("{state_path}.db")); + let _ = std::fs::remove_file(format!("{state_path}.audit.jsonl")); + let _ = std::fs::remove_dir_all(sessions_dir); + let _ = std::fs::remove_dir_all(assistant_dir); } diff --git a/src/gateway/worker.rs b/src/gateway/worker.rs index ac7a0d5..c366a62 100644 --- a/src/gateway/worker.rs +++ b/src/gateway/worker.rs @@ -412,7 +412,11 @@ where ctx.audit .backend_completed(job.row_id, &job.thread, job.backend, &out.reply), ); - let (reply_quote, reply_text) = super::lift_outbound_quote(&out.reply); + let (reply_quote, reply_text) = if ctx.channel.supports_reply_quotes() { + super::lift_outbound_quote(&out.reply) + } else { + (None, out.reply.clone()) + }; let outbound = match ctx.history.lock().unwrap().record_outbound( job.inbound_id, OutboundOrigin::Backend,