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/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/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"]; diff --git a/src/audit.rs b/src/audit.rs index eb87454..34132a9 100644 --- a/src/audit.rs +++ b/src/audit.rs @@ -395,6 +395,8 @@ mod tests { images: Vec::new(), 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 4dc8720..4ffebd1 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -55,6 +55,12 @@ 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, + /// 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 { @@ -70,6 +76,11 @@ 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, + /// 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)] @@ -214,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, @@ -422,6 +439,8 @@ impl ChannelContract for IMessageChannel { .collect(), is_from_me: message.is_from_me, is_supported: true, + reply_to_message_id: None, + reply_context: None, thread_id: None, }) .collect()) @@ -515,6 +534,8 @@ impl ChannelContract for IMessageChannel { Vec::new() } else { vec![OutboundChunk { + reply_to_message_id: None, + quote: None, text: format!("{text}{marker}"), rich_markdown: false, }] @@ -630,6 +651,8 @@ impl ChannelContract for Telegram { crate::telegram::split_text(text) .into_iter() .map(|text| OutboundChunk { + reply_to_message_id: None, + quote: None, text, rich_markdown: true, }) @@ -638,9 +661,21 @@ 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, + chunk.quote.as_deref(), + ) + .await } else { - self.send_plain(target, &chunk.text).await + self.send_plain_reply( + target, + &chunk.text, + chunk.reply_to_message_id, + chunk.quote.as_deref(), + ) + .await } } @@ -736,6 +771,8 @@ 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, + quote: None, text, rich_markdown: true, }) @@ -883,6 +920,8 @@ mod tests { images: Vec::new(), is_from_me, is_supported: true, + reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -900,6 +939,8 @@ mod tests { images: Vec::new(), is_from_me: false, is_supported: true, + reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -939,6 +980,8 @@ mod tests { assert_eq!( channel.outbound_chunks("reply", REPLY_MARKER), [OutboundChunk { + reply_to_message_id: None, + quote: None, text: format!("reply{REPLY_MARKER}"), rich_markdown: false, }] @@ -974,6 +1017,8 @@ mod tests { assert_eq!( channel.outbound_chunks("reply", REPLY_MARKER), [OutboundChunk { + reply_to_message_id: None, + quote: None, text: "reply".to_string(), rich_markdown: true, }] @@ -995,6 +1040,8 @@ mod tests { images: Vec::new(), 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 e02b687..efc6cde 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. @@ -208,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, @@ -455,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 { @@ -787,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"); } } @@ -824,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 { @@ -889,20 +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; } @@ -1224,14 +1237,38 @@ fn audit_schedule_events(ctx: &Ctx, ledger: &mut jobs::Ledger) { } } -async fn reply_to(ctx: &Ctx, target: &str, text: &str) -> bool { +/// 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() { 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, &chunk).await { + if let Err(error) = send_reply_chunk(ctx, target, reply_anchor, None, &chunk).await { error!("send error to {target}: {error}"); return false; } @@ -1239,9 +1276,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, + reply_anchor: Option, + quote: Option<&str>, chunk: &crate::channel::OutboundChunk, ) -> Result<()> { #[cfg(test)] @@ -1277,11 +1317,14 @@ async fn send_reply_chunk( } #[cfg(not(test))] { + 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 + 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"), } @@ -1325,7 +1368,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, 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 c7a0bd0..1a1888e 100644 --- a/src/gateway/tests.rs +++ b/src/gateway/tests.rs @@ -166,6 +166,8 @@ 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, + reply_context: None, thread_id: None, } } @@ -239,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(), @@ -396,7 +399,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())); @@ -1074,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); @@ -1158,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( @@ -1170,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); @@ -1884,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, }), )])); @@ -1943,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, ""); @@ -2473,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, }), )])); @@ -2633,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( @@ -2786,6 +2796,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(), @@ -2872,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); @@ -2981,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); @@ -3096,6 +3109,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(), @@ -3248,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); @@ -3313,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); @@ -4035,6 +4051,7 @@ fn fake_runners_with_hook( wait_for_release: None, failure: None, resume_missing_once: None, + reply: None, }), ); runners @@ -4096,6 +4113,8 @@ 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, + reply_context: None, thread_id: None, } } @@ -4119,6 +4138,8 @@ fn telegram_message( is_from_me: false, is_group, is_supported: true, + reply_to_message_id: None, + reply_context: None, thread_id: None, } } @@ -4173,6 +4194,135 @@ fn slack_image_message( is_from_me: false, is_group: false, is_supported: true, + reply_to_message_id: None, + reply_context: None, 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"); +} + +#[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"); + 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; + + // 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(), + "> 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 ecba8df..c366a62 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", @@ -412,11 +412,16 @@ where ctx.audit .backend_completed(job.row_id, &job.thread, job.backend, &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, Some(job.backend.as_str()), - &out.reply, + &reply_text, ) { Ok(outbound) => outbound, Err(error) => { @@ -443,7 +448,7 @@ where ); return; } - let delivery = deliver_stored(ctx, &job, &outbound).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 +456,7 @@ where ctx, &job, delivery, - &out.reply, + &reply_text, "completed", "deliver backend reply", ); @@ -576,8 +581,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 { @@ -713,7 +723,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, @@ -848,13 +858,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 @@ -874,7 +885,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, &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 { @@ -890,6 +908,7 @@ async fn deliver_stored( ctx: &Ctx, job: &Job, outbound: &OutboundMessage, + quote: Option<&str>, ) -> Result { if outbound.status == DeliveryStatus::Delivered { return Ok(DeliveryOutcome::AlreadyDelivered); @@ -899,7 +918,14 @@ 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, + quote, + &mut outbound, + ) + .await?; let status = if delivered { DeliveryStatus::Delivered } else { @@ -944,6 +970,8 @@ async fn deliver_stored( async fn deliver_outbound_once( ctx: &Ctx, target: &str, + reply_anchor: Option, + quote: Option<&str>, outbound: &mut OutboundMessage, ) -> Result { let chunks = ctx @@ -962,7 +990,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, reply_anchor, quote, chunk).await { error!( "outbound {} chunk {index} send error to {target}: {error}", outbound.id diff --git a/src/slack.rs b/src/slack.rs index 8e47080..778945c 100644 --- a/src/slack.rs +++ b/src/slack.rs @@ -663,6 +663,8 @@ impl Inbox { .collect(), 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 2505e34..d104f9a 100644 --- a/src/telegram.rs +++ b/src/telegram.rs @@ -262,7 +262,13 @@ 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<()> { + pub async fn send_rich_reply( + &self, + 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"); } @@ -270,6 +276,9 @@ impl Telegram { let mut payload = target_payload(target); payload["text"] = json!(html); payload["parse_mode"] = json!("HTML"); + 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) .await?; @@ -278,15 +287,26 @@ 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<()> { + pub async fn send_plain_reply( + &self, + 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(params) = reply_parameters_payload(reply_to_message_id, quote) { + payload["reply_parameters"] = params; + } let transport_response = self .post_with_topic_fallback("sendMessage", payload) .await?; @@ -501,6 +521,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 { @@ -517,6 +554,8 @@ impl Update { is_from_me: false, is_supported: false, thread_id: None, + reply_to_message_id: None, + reply_context: None, }; }; let images = message @@ -572,16 +611,41 @@ impl Update { is_from_me: false, is_supported: true, thread_id: message.message_thread_id, + reply_to_message_id: message.message_id, + reply_context: message + .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)), } } } +/// 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)] from: Option, chat: Chat, #[serde(default)] + message_id: Option, + #[serde(default)] text: Option, #[serde(default)] caption: Option, @@ -593,6 +657,16 @@ struct TelegramMessage { voice: Option, #[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)] @@ -819,6 +893,108 @@ 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(); + 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}}); + 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_reply("chat", "hi", None, None) + .await + .unwrap(); + telegram + .send_plain_reply("chat", "hi", Some(41), Some("the frobnicator")) + .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 + ); + 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] async fn poll_uses_next_update_offset_and_long_poll_timeout() { let fake = Arc::new(FakeTransport::with_responses(vec![json!({ @@ -867,6 +1043,9 @@ mod tests { document: None, voice: None, message_thread_id: None, + message_id: Some(1), + quote: None, + reply_to_message: None, }), } .into_raw(); @@ -898,6 +1077,9 @@ mod tests { mime_type: Some("audio/ogg".to_string()), }), message_thread_id: None, + message_id: Some(1), + quote: None, + reply_to_message: None, }), } .into_raw(); @@ -1087,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!( @@ -1114,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. @@ -1132,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(); @@ -1150,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(); @@ -1190,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!( @@ -1217,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); @@ -1229,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"); @@ -1286,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();