Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions src/approvals.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@

use std::sync::Arc;

use crate::platform::{Button, OutboundMessage, Platform, PlatformResult, ReplyCtx};
use crate::platform::{Button, MessageRef, OutboundMessage, Platform, PlatformResult, ReplyCtx};
use crate::render::PendingAsk;

/// Longest a button's visible label is allowed to be — Telegram (and most IM
Expand Down Expand Up @@ -116,15 +116,16 @@ pub fn format_ask_text(ask: &ParkedAsk) -> String {
/// Render `ask` to the chat: inline buttons (one per option, `callback_data =
/// "{nonce}:{index}"`) when `buttons_capable`, always alongside the numbered
/// text list — per bamboo issue #458, buttons are an enhancement, never a
/// requirement. Returns the platform error (if the send failed) so the
/// caller can log it; rendering failure does not itself invalidate the
/// parked ask (a text reply can still resolve it).
/// requirement. Returns the sent message's [`MessageRef`] so the bridge can
/// edit the ask once answered (✅ + chosen answer, buttons dropped — issue
/// #6 follow-up); a send failure is returned for logging but does not itself
/// invalidate the parked ask (a text reply can still resolve it).
pub async fn render_ask(
platform: &Arc<dyn Platform>,
reply_ctx: &ReplyCtx,
ask: &ParkedAsk,
buttons_capable: bool,
) -> PlatformResult<()> {
) -> PlatformResult<MessageRef> {
let text = format_ask_text(ask);
let outbound = if buttons_capable && !ask.options.is_empty() {
let rows: Vec<Vec<Button>> = ask
Expand All @@ -142,7 +143,7 @@ pub async fn render_ask(
} else {
OutboundMessage::text(text)
};
platform.reply(reply_ctx, outbound).await.map(|_| ())
platform.reply(reply_ctx, outbound).await
}

// ---------------------------------------------------------------------------
Expand Down
102 changes: 91 additions & 11 deletions src/bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -895,12 +895,20 @@ impl ConnectBridge {
let caps = platform.capabilities();
let parked =
ParkedAsk::new(approvals::new_nonce(), session_id.to_string(), &ask);

if let Err(error) =
approvals::render_ask(&platform, &reply_ctx, &parked, caps.buttons).await
{
tracing::warn!("magpie bridge: failed to render pending ask: {error}");
}
let ask_text = approvals::format_ask_text(&parked);

let ask_ref =
match approvals::render_ask(&platform, &reply_ctx, &parked, caps.buttons)
.await
{
Ok(msg_ref) => Some(msg_ref),
Err(error) => {
tracing::warn!(
"magpie bridge: failed to render pending ask: {error}"
);
None
}
};

let (ask_tx, mut ask_rx) = mpsc::channel(1);
{
Expand All @@ -912,13 +920,51 @@ impl ConnectBridge {

match ask_rx.recv().await {
Some(AskResolution::Answer(answer)) => {
// Subscribe BEFORE respond — the same
// subscribe-before-execute invariant the initial
// run relies on (ARCHITECTURE.md). A server-side
// resubscribe REPLACES the channel's forwarder
// with a fresh broadcast cut, so subscribing
// after `respond` leaves a window where the
// resumed run's events — or, if the WS is mid-
// reconnect while the Subscribe command waits,
// the entire run through `Complete` — are
// emitted into no subscription at all; the
// render task then waits forever on an idle
// channel and the final reply never renders
// (issue #6). A failed subscribe still records
// the answer below (the old degradation path):
// dropping the user's decision is strictly worse
// than rendering nothing.
let new_rx = self.api.subscribe_session(session_id).await;
let respond_request = RespondRequest {
response: answer,
response: answer.clone(),
..Default::default()
};
match self.api.respond(session_id, respond_request).await {
Ok(_response) => {
match self.api.subscribe_session(session_id).await {
// Mark the ask message answered — ✅ +
// the chosen answer, buttons dropped (an
// edit replaces the whole message body) —
// so stale buttons can't be pressed again
// and the chat shows WHAT was chosen.
// Best-effort: an edit failure never
// fails the resume.
if caps.edit_message {
if let Some(msg_ref) = &ask_ref {
let done = format!("{ask_text}\n\n✅ {answer}");
if let Err(error) = platform
.edit(msg_ref, OutboundMessage::text(done))
.await
{
tracing::debug!(
"magpie bridge: answered-ask edit failed \
(non-fatal): {error}"
);
}
}
}
match new_rx {
Ok(new_rx) => {
rx = new_rx;
continue;
Expand Down Expand Up @@ -1090,6 +1136,11 @@ mod tests {
respond_calls: TokioMutex<Vec<(String, String)>>,
respond_pending_calls: TokioMutex<Vec<String>>,
subscribe_calls: TokioMutex<Vec<String>>,
/// Interleaved `"subscribe"`/`"respond"` markers in true call order —
/// the per-method vecs above can't express cross-method ordering,
/// which the resume path's subscribe-BEFORE-respond invariant
/// (issue #6) needs asserted.
ops: TokioMutex<Vec<&'static str>>,
channels: TokioMutex<HashMap<String, mpsc::Sender<StreamEvent>>>,
next_id: AtomicUsize,
chat_error: Option<String>,
Expand All @@ -1105,6 +1156,7 @@ mod tests {
respond_calls: TokioMutex::new(Vec::new()),
respond_pending_calls: TokioMutex::new(Vec::new()),
subscribe_calls: TokioMutex::new(Vec::new()),
ops: TokioMutex::new(Vec::new()),
channels: TokioMutex::new(HashMap::new()),
next_id: AtomicUsize::new(1),
chat_error: None,
Expand All @@ -1120,6 +1172,7 @@ mod tests {
respond_calls: TokioMutex::new(Vec::new()),
respond_pending_calls: TokioMutex::new(Vec::new()),
subscribe_calls: TokioMutex::new(Vec::new()),
ops: TokioMutex::new(Vec::new()),
channels: TokioMutex::new(HashMap::new()),
next_id: AtomicUsize::new(1),
chat_error: None,
Expand Down Expand Up @@ -1196,6 +1249,7 @@ mod tests {
.lock()
.await
.push((session_id.to_string(), request.response.clone()));
self.ops.lock().await.push("respond");
if let Some(message) = &self.respond_error {
return Err(ClientError::Api {
method: "POST",
Expand Down Expand Up @@ -1240,6 +1294,7 @@ mod tests {
.lock()
.await
.push(session_id.to_string());
self.ops.lock().await.push("subscribe");
let (tx, rx) = mpsc::channel(16);
self.channels
.lock()
Expand Down Expand Up @@ -1621,7 +1676,7 @@ mod tests {
fn buttons_and_edit_capabilities() -> crate::platform::Capabilities {
crate::platform::Capabilities {
buttons: true,
edit_message: false,
edit_message: true,
images: false,
files: false,
}
Expand Down Expand Up @@ -1706,13 +1761,38 @@ mod tests {
);
assert_eq!(platform.answered_callbacks.lock().await.len(), 1);

// Bridge re-subscribed after respond (per ARCHITECTURE.md) — end the
// resumed run so the background task settles.
// The resume leg re-subscribes BEFORE responding (issue #6): a
// server-side resubscribe replaces the channel's forwarder with a
// fresh broadcast cut, so subscribing after `respond` can lose the
// resumed run's events entirely — only the ordered op log can see
// the cross-method ordering.
wait_until(|| {
let api = api.clone();
async move { api.subscribe_calls.lock().await.len() == 2 }
})
.await;
assert_eq!(
api.ops.lock().await.as_slice(),
["subscribe", "subscribe", "respond"],
"resume must subscribe before respond"
);

// The answered ask message was edited to ✅ + the chosen answer, so
// its now-stale buttons can't be pressed again (issue #6 follow-up).
wait_until(|| {
let platform = platform.clone();
async move {
platform
.edits
.lock()
.await
.iter()
.any(|edit| edit.contains("✅ Approve"))
}
})
.await;

// End the resumed run so the background task settles.
let session_id = bridge.session_id_for_key(&key).await.unwrap();
api.sender_for(&session_id)
.await
Expand Down
18 changes: 16 additions & 2 deletions src/render.rs
Original file line number Diff line number Diff line change
Expand Up @@ -356,11 +356,18 @@ impl StreamingRenderer {
fn resume(platform: Arc<dyn Platform>, reply_ctx: ReplyCtx, state: StreamState) -> Self {
let StreamState {
tool_lines,
assistant_text,
mut assistant_text,
status_ref,
last_edit_at,
chars_since_edit,
} = state;
// Paragraph break between the pre-pause text and whatever the
// resumed run streams next — without it the resumed reply glues
// straight onto the question ("Pick one:OASIS", issue #6). Trailing
// whitespace is trimmed at finalize if no tokens follow.
if !assistant_text.is_empty() && !assistant_text.ends_with('\n') {
assistant_text.push_str("\n\n");
}
Self {
platform,
reply_ctx,
Expand Down Expand Up @@ -484,7 +491,7 @@ impl StreamingRenderer {
/// edit plus the full result sent as fresh chunked messages (bamboo issue
/// #458 §B point 5).
async fn finalize_success(&mut self) {
let full = self.full_body();
let full = self.full_body().trim_end().to_string();
if full.trim().is_empty() {
self.apply_edit("✅ Done.".to_string()).await;
return;
Expand Down Expand Up @@ -1039,6 +1046,13 @@ mod tests {
assert!(last.starts_with('✅'), "final edit not a success: {last}");
assert!(last.contains("Before the question."));
assert!(last.contains("After the answer."));
// The resumed segment starts on its own paragraph — without the
// separator the resumed reply glues straight onto the pre-pause text
// ("Pick one:OASIS", issue #6).
assert!(
last.contains("Before the question. \n\nAfter the answer."),
"resumed text must be separated from pre-pause text: {last}"
);
}

#[tokio::test]
Expand Down
Loading