From a6f0fc0c0f59e05faf1b8621cec9fb6d97e31ca8 Mon Sep 17 00:00:00 2001 From: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> Date: Tue, 11 Aug 2026 12:38:21 -0400 Subject: [PATCH] fix(acp): release prepared edits before fallback Signed-off-by: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> --- crates/buzz-acp/src/lib.rs | 6 +++- crates/buzz-acp/src/queue.rs | 65 ++++++++++++++++++++++++++++-------- 2 files changed, 56 insertions(+), 15 deletions(-) diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 1e03b8cc0..122da7165 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -2627,9 +2627,13 @@ async fn tokio_main() -> Result<()> { if queue.has_native_steer_reservations( buzz_event.channel_id, ) { + let released = queue.release_native_steers( + buzz_event.channel_id, + ); tracing::debug!( channel_id = %buzz_event.channel_id, - "native steer preparation already pending; preserving channel order via cancel+merge" + released, + "native steer preparation already pending; releasing reservations for ordered cancel+merge" ); false } else if is_edit diff --git a/crates/buzz-acp/src/queue.rs b/crates/buzz-acp/src/queue.rs index 84c77bc24..3d03b6086 100644 --- a/crates/buzz-acp/src/queue.rs +++ b/crates/buzz-acp/src/queue.rs @@ -796,6 +796,29 @@ impl EventQueue { } } + /// Release every native-steer reservation for `channel_id` to the queue + /// front, preserving original FIFO order. Detached edit preparation then + /// observes its missing reservation and discards the stale completion. + pub fn release_native_steers(&mut self, channel_id: Uuid) -> usize { + let Some(entries) = self.withheld_native_steer.remove(&channel_id) else { + return 0; + }; + let n = entries.len(); + let queue = self.queues.entry(channel_id).or_default(); + for qe in entries.into_iter().rev() { + queue.push_front(qe); + } + while queue.len() > MAX_PENDING_PER_CHANNEL { + queue.pop_back(); + tracing::warn!( + channel_id = %channel_id, + limit = MAX_PENDING_PER_CHANNEL, + "withheld-steer release overflow — dropped newest event to enforce cap" + ); + } + n + } + /// Bulk-release every withheld event for `channel_id` back to the queue /// front, preserving relative FIFO order. /// @@ -810,21 +833,9 @@ impl EventQueue { /// composes to original-FIFO order at the queue front (same discipline /// as `requeue_preserve_timestamps` at line 453). fn recover_withheld_for_expired_channel(&mut self, channel_id: Uuid) { - let Some(entries) = self.withheld_native_steer.remove(&channel_id) else { + let n = self.release_native_steers(channel_id); + if n == 0 { return; - }; - let n = entries.len(); - let queue = self.queues.entry(channel_id).or_default(); - for qe in entries.into_iter().rev() { - queue.push_front(qe); - } - while queue.len() > MAX_PENDING_PER_CHANNEL { - queue.pop_back(); - tracing::warn!( - channel_id = %channel_id, - limit = MAX_PENDING_PER_CHANNEL, - "withheld-steer recovery overflow — dropped newest event to enforce cap" - ); } tracing::warn!( channel_id = %channel_id, @@ -4758,6 +4769,32 @@ mod tests { )); } + #[test] + fn releasing_pending_preparation_restores_fifo_ahead_of_later_event() { + let mut q = EventQueue::new(DedupMode::Queue); + let ch = Uuid::new_v4(); + let edit = make_queued_at(ch, "edit", Duration::from_millis(20)); + let edit_id = edit.event.id.to_hex(); + let later = make_queued_at(ch, "later", Duration::from_millis(10)); + let later_id = later.event.id.to_hex(); + q.push(edit); + assert!(q.mark_native_steer_pending(ch, &edit_id)); + q.push(later); + + // A later event triggers cancel+merge. Releasing the pending edit first + // makes both events visible to the next flush in original arrival order; + // the detached preparation will be stale because its reservation is gone. + assert_eq!(q.release_native_steers(ch), 1); + assert!(!q.has_native_steer_reservation(ch, &edit_id)); + let replacement = q.flush_next().expect("both events should dispatch"); + let ids: Vec<_> = replacement + .events + .iter() + .map(|event| event.event.id.to_hex()) + .collect(); + assert_eq!(ids, vec![edit_id, later_id]); + } + /// Bulk-release on expiry must preserve original FIFO. The /// implementation iterates the side-table entries in reverse and /// `push_front`s each — composing to original-FIFO at the queue front.