Skip to content
Open
6 changes: 6 additions & 0 deletions assistant/skills/push/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `> <that part>` 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.
12 changes: 9 additions & 3 deletions src/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -177,6 +177,9 @@ pub struct FakeRunner {
pub wait_for_release: Option<std::sync::Arc<tokio::sync::Notify>>,
pub failure: Option<String>,
pub resume_missing_once: Option<std::sync::Arc<std::sync::atomic::AtomicBool>>,
/// Overrides the canned reply so tests can exercise backend-reply
/// formatting conventions (for example a leading `> quoted portion`).
pub reply: Option<String>,
}

#[cfg(test)]
Expand Down Expand Up @@ -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()),
})
}
Expand Down
2 changes: 1 addition & 1 deletion src/assistant.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"];
Expand Down
2 changes: 2 additions & 0 deletions src/audit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
}
Expand Down
51 changes: 49 additions & 2 deletions src/channel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,12 @@ pub struct RawMessage {
pub is_supported: bool,
/// Channel-specific thread/topic id (Telegram `message_thread_id`).
pub thread_id: Option<i64>,
/// Provider message id of the inbound message (Telegram `message_id`),
/// used to anchor replies via `reply_parameters`.
pub reply_to_message_id: Option<i64>,
/// Flattened excerpt of the message the user replied to (Telegram
/// `reply_to_message`), prefixed to the backend prompt. None elsewhere.
pub reply_context: Option<String>,
}

impl RawMessage {
Expand All @@ -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<i64>,
/// Exact portion of the trigger message to quote in the reply header
/// (Telegram `reply_parameters.quote`), chosen by the backend.
pub quote: Option<String>,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
Expand Down Expand Up @@ -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<Vec<RawMessage>> {
match self {
Self::IMessage(channel) => ChannelContract::poll(channel, since).await,
Expand Down Expand Up @@ -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())
Expand Down Expand Up @@ -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,
}]
Expand Down Expand Up @@ -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,
})
Expand All @@ -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
}
}

Expand Down Expand Up @@ -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,
})
Expand Down Expand Up @@ -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,
}
}
Expand All @@ -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,
}
}
Expand Down Expand Up @@ -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,
}]
Expand Down Expand Up @@ -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,
}]
Expand All @@ -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,
};

Expand Down
65 changes: 54 additions & 11 deletions src/gateway/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,8 @@ struct Job {
voice_attachment: Option<InboundVoice>,
image_attachments: Vec<InboundImage>,
approval_origin: AnswerOrigin,
/// Telegram message id of the trigger message (reply anchoring).
telegram_reply_anchor: Option<i64>,
}

/// Shared, cheaply cloneable context handed to each worker task.
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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");
}
}
Expand Down Expand Up @@ -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
{
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -1224,24 +1237,51 @@ 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>, 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<i64>) -> 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;
}
}
true
}

#[cfg_attr(test, allow(unused_variables))]
async fn send_reply_chunk(
ctx: &Ctx,
target: &str,
reply_anchor: Option<i64>,
quote: Option<&str>,
chunk: &crate::channel::OutboundChunk,
) -> Result<()> {
#[cfg(test)]
Expand Down Expand Up @@ -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"),
}
Expand Down Expand Up @@ -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<Mutex<Store>>, ack: &Arc<Mutex<AckState>>, channel: &str, row_id: i64) {
Expand Down
Loading