diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 2a339e101..7c456aa5b 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -2554,6 +2554,12 @@ async fn tokio_main() -> Result<()> { let native_attempted = if matches!(signal, ControlSignal::Steer) { if queue::edit_target_id(&event_for_steer).is_some() { + let event_id = event_for_steer.id.to_hex(); + let reserved = queue.mark_native_steer_pending( + buzz_event.channel_id, + &event_id, + ); + debug_assert!(reserved, "accepted edit must still be queued"); let tx = native_steer_tx.clone(); let ctx = Arc::clone(&ctx); let channel_id = buzz_event.channel_id; @@ -2790,10 +2796,11 @@ async fn tokio_main() -> Result<()> { &mut pool, &mut queue, channel_id, - event, + event.clone(), prompt_blocks, &steer_ack_tx, ) { + queue.release_native_steer(channel_id, &event.id.to_hex()); signal_in_flight_task(&mut pool, channel_id, ControlSignal::Steer); } } @@ -3237,25 +3244,14 @@ fn try_native_steer( match pool.send_steer(channel_id, request) { Ok(()) => { - // Withhold the queued event synchronously BEFORE spawning - // the watcher: this closes the race where `mark_complete` - // clears `in_flight_channels` and a stray `flush_next` could - // re-deliver the event via normal dispatch. See - // `EventQueue::mark_native_steer_pending` docs at queue.rs:606. + // Ordinary events are withheld after send. Edit preparation reserves + // its event before leaving the main loop, so this is idempotent. let withheld = queue.mark_native_steer_pending(channel_id, &event_id_hex); if !withheld { - // Race: the event was already drained out of the queue - // before we got here (e.g. a concurrent flush picked it - // up). The steer is on the wire; if it succeeds the - // agent gets it via the native path AND normal - // dispatch — duplicate delivery is benign (agent gets - // the same message twice). Log so this is visible if it - // ever happens in production. - tracing::warn!( + tracing::debug!( channel = %channel_id, event_id = %event_id_hex, - "native steer accepted by read loop but event was not in queue to withhold \ - — possible duplicate delivery if steer succeeds" + "native steer event was already reserved during async preparation" ); } let ack_tx_clone = steer_ack_tx.clone(); diff --git a/crates/buzz-acp/src/queue.rs b/crates/buzz-acp/src/queue.rs index 459351f5c..e3ee95370 100644 --- a/crates/buzz-acp/src/queue.rs +++ b/crates/buzz-acp/src/queue.rs @@ -5071,6 +5071,34 @@ mod tests { assert!(!prompt.contains(&format!("--reply-to {edit_id}"))); } + #[test] + fn async_edit_steer_reservation_prevents_normal_redispatch() { + let mut q = EventQueue::new(DedupMode::Queue); + let ch = Uuid::new_v4(); + let event = edit_event(&"aa".repeat(32)); + let event_id = event.id.to_hex(); + q.push(QueuedEvent { + channel_id: ch, + event, + received_at: Instant::now(), + prompt_tag: "@mention".into(), + }); + q.in_flight_channels.insert(ch); + q.in_flight_deadlines + .insert(ch, Instant::now() + Duration::from_secs(60)); + + assert!(q.mark_native_steer_pending(ch, &event_id)); + q.mark_complete(ch); + assert!( + q.flush_next().is_none(), + "reserved edit must not redispatch" + ); + + q.release_native_steer(ch, &event_id); + let batch = q.flush_next().expect("failed preparation restores edit"); + assert_eq!(batch.events[0].event.id.to_hex(), event_id); + } + #[test] fn native_steer_edit_uses_original_thread_anchor() { let original_id = "66".repeat(32);