mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(acp): reserve edits during steer preparation
Signed-off-by: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz>
This commit is contained in:
+12
-16
@@ -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();
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user