From 3e668cdfcfc9b5a5875e856741be5483f41eb576 Mon Sep 17 00:00:00 2001 From: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> Date: Mon, 10 Aug 2026 14:11:51 -0400 Subject: [PATCH] fix(acp): reserve edits during steer preparation Signed-off-by: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> --- crates/buzz-acp/src/lib.rs | 28 ++++++++++++---------------- crates/buzz-acp/src/queue.rs | 28 ++++++++++++++++++++++++++++ 2 files changed, 40 insertions(+), 16 deletions(-) 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);