mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(acp): complete edit lifecycle routing
Use resolved original-thread routing for edit-triggered typing and return visible reaction targets when draining queued or reserved edits.
Signed-off-by: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz>
(cherry picked from commit 05675c3129)
Signed-off-by: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz>
This commit is contained in:
+96
-20
@@ -2546,7 +2546,9 @@ async fn tokio_main() -> Result<()> {
|
||||
&ctx,
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
) {
|
||||
)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -2599,7 +2601,9 @@ async fn tokio_main() -> Result<()> {
|
||||
&ctx,
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
) {
|
||||
)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -3191,6 +3195,7 @@ async fn tokio_main() -> Result<()> {
|
||||
if pool_ready {
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &membership_generations, &mut last_activity)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
@@ -3291,6 +3296,7 @@ async fn tokio_main() -> Result<()> {
|
||||
tracing::debug!("heartbeat_skipped_events");
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &membership_generations, &mut last_activity)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
@@ -3396,7 +3402,9 @@ async fn tokio_main() -> Result<()> {
|
||||
&ctx,
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
) {
|
||||
)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -3426,7 +3434,9 @@ async fn tokio_main() -> Result<()> {
|
||||
&ctx,
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
) {
|
||||
)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -3477,7 +3487,8 @@ async fn tokio_main() -> Result<()> {
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
&mut typing_channels,
|
||||
);
|
||||
)
|
||||
.await;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
@@ -3499,7 +3510,8 @@ async fn tokio_main() -> Result<()> {
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
&mut typing_channels,
|
||||
);
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
Some(PoolEvent::SteerAck(SteerAckEvent {
|
||||
@@ -3682,7 +3694,9 @@ async fn tokio_main() -> Result<()> {
|
||||
&ctx,
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
) {
|
||||
)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -3714,7 +3728,9 @@ async fn tokio_main() -> Result<()> {
|
||||
&ctx,
|
||||
&membership_generations,
|
||||
&mut last_activity,
|
||||
) {
|
||||
)
|
||||
.await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -4049,7 +4065,7 @@ fn try_native_steer(
|
||||
/// universal cancel+merge fallback, and immediately try dispatch in case the
|
||||
/// original turn already ended.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn release_prepared_native_steer_fallback(
|
||||
async fn release_prepared_native_steer_fallback(
|
||||
pool: &mut AgentPool,
|
||||
queue: &mut EventQueue,
|
||||
channel_id: Uuid,
|
||||
@@ -4071,7 +4087,7 @@ fn release_prepared_native_steer_fallback(
|
||||
ControlSignal::Steer,
|
||||
);
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(pool, queue, ctx, membership_generations, last_activity)
|
||||
dispatch_pending(pool, queue, ctx, membership_generations, last_activity).await
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
@@ -4079,8 +4095,24 @@ fn release_prepared_native_steer_fallback(
|
||||
|
||||
// ── dispatch_pending ──────────────────────────────────────────────────────────
|
||||
|
||||
fn typing_scope_for_batch(
|
||||
batch: &FlushBatch,
|
||||
resolved_edit: Option<&queue::ResolvedEdit>,
|
||||
) -> ThreadTags {
|
||||
resolved_edit
|
||||
.map(queue::ResolvedEdit::reply_thread_tags)
|
||||
.or_else(|| batch.failure_thread_tags.clone())
|
||||
.or_else(|| {
|
||||
batch
|
||||
.events
|
||||
.last()
|
||||
.map(|event| queue::parse_thread_tags(&event.event))
|
||||
})
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
/// Flush queued work to available agents.
|
||||
fn dispatch_pending(
|
||||
async fn dispatch_pending(
|
||||
pool: &mut AgentPool,
|
||||
queue: &mut EventQueue,
|
||||
ctx: &Arc<PromptContext>,
|
||||
@@ -4094,11 +4126,15 @@ fn dispatch_pending(
|
||||
None => break,
|
||||
};
|
||||
let channel_id = batch.channel_id;
|
||||
let typing_scope = batch
|
||||
.events
|
||||
.last()
|
||||
.map(|event| queue::parse_thread_tags(&event.event))
|
||||
.unwrap_or_default();
|
||||
// Typing indicators are user-visible routing and must share the same
|
||||
// resolved-edit authority as prompts, replies, reactions, and failure
|
||||
// notices. Resolve before spawning the prompt task so the main loop
|
||||
// never briefly advertises an edit-triggered turn at channel scope.
|
||||
let resolved_typing_edit = match batch.events.last() {
|
||||
Some(event) => pool::resolve_edit_routing(&event.event, &ctx.rest_client).await,
|
||||
None => None,
|
||||
};
|
||||
let typing_scope = typing_scope_for_batch(&batch, resolved_typing_edit.as_ref());
|
||||
let affinity_hit = pool.has_session_for(channel_id);
|
||||
let mut agent = match pool.try_claim(Some(channel_id)) {
|
||||
Some(a) => a,
|
||||
@@ -9301,6 +9337,41 @@ mod native_edit_membership_lifecycle_tests {
|
||||
assert_eq!(actual.parent_event_id.as_deref(), Some(parent.as_str()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn edit_typing_scope_uses_the_resolved_original_thread() {
|
||||
let channel_id = Uuid::new_v4();
|
||||
let target = "ab".repeat(32);
|
||||
let root = "cd".repeat(32);
|
||||
let parent = "ef".repeat(32);
|
||||
let batch = FlushBatch {
|
||||
channel_id,
|
||||
events: vec![queue::BatchEvent {
|
||||
event: edit_event(&target),
|
||||
prompt_tag: "@mention".into(),
|
||||
received_at: std::time::Instant::now(),
|
||||
}],
|
||||
cancelled_events: vec![],
|
||||
cancel_reason: None,
|
||||
failure_thread_tags: Some(ThreadTags {
|
||||
root_event_id: Some(target.clone()),
|
||||
parent_event_id: Some(target),
|
||||
mentioned_pubkeys: vec![],
|
||||
}),
|
||||
};
|
||||
let resolved = queue::ResolvedEdit {
|
||||
target_event_id: "12".repeat(32),
|
||||
target_thread_tags: ThreadTags {
|
||||
root_event_id: Some(root.clone()),
|
||||
parent_event_id: Some(parent.clone()),
|
||||
mentioned_pubkeys: vec![],
|
||||
},
|
||||
};
|
||||
|
||||
let actual = typing_scope_for_batch(&batch, Some(&resolved));
|
||||
assert_eq!(actual.root_event_id.as_deref(), Some(root.as_str()));
|
||||
assert_eq!(actual.parent_event_id.as_deref(), Some(parent.as_str()));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn prepared_native_steer_rejection_releases_and_dispatches_immediately() {
|
||||
let channel_id = Uuid::new_v4();
|
||||
@@ -9330,7 +9401,8 @@ mod native_edit_membership_lifecycle_tests {
|
||||
&HashMap::new(),
|
||||
&mut last_activity,
|
||||
&mut typing_channels,
|
||||
);
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(typing_channels.contains_key(&channel_id));
|
||||
assert_eq!(
|
||||
@@ -9444,13 +9516,17 @@ mod native_edit_membership_lifecycle_tests {
|
||||
// drains the reservation and advances generation; re-add restores
|
||||
// access without revalidating old prepared work.
|
||||
let mut removed_channels = HashSet::new();
|
||||
assert!(remove_channel_for_membership_change(
|
||||
let drained = remove_channel_for_membership_change(
|
||||
&mut queue,
|
||||
&mut removed_channels,
|
||||
&mut generations,
|
||||
channel_id,
|
||||
)
|
||||
.is_empty());
|
||||
);
|
||||
assert_eq!(
|
||||
drained,
|
||||
vec!["ab".repeat(32)],
|
||||
"removal returns the reserved edit's visible reaction target"
|
||||
);
|
||||
readd_channel_after_membership_change(&mut removed_channels, channel_id);
|
||||
|
||||
// This is the completion-arm guard. A false result would enter
|
||||
|
||||
@@ -770,19 +770,26 @@ impl EventQueue {
|
||||
///
|
||||
/// Also clears any `retry_after` throttle for the channel.
|
||||
///
|
||||
/// Returns the event IDs of dropped events so the caller can clean up
|
||||
/// any reactions (👀) that were added at queue-push time.
|
||||
/// Returns the visible event IDs that own lifecycle reactions for every
|
||||
/// dropped queued or reserved event. Edit events therefore contribute
|
||||
/// their original target IDs rather than their auxiliary kind:40003 IDs.
|
||||
pub fn drain_channel(&mut self, channel_id: Uuid) -> Vec<String> {
|
||||
let ids = self
|
||||
let mut ids: Vec<String> = self
|
||||
.queues
|
||||
.remove(&channel_id)
|
||||
.map(|q| q.into_iter().map(|e| e.event.id.to_hex()).collect())
|
||||
.map(|q| {
|
||||
q.into_iter()
|
||||
.map(|e| reaction_target_id(&e.event))
|
||||
.collect()
|
||||
})
|
||||
.unwrap_or_default();
|
||||
if let Some(withheld) = self.withheld_native_steer.remove(&channel_id) {
|
||||
ids.extend(withheld.into_iter().map(|e| reaction_target_id(&e.event)));
|
||||
}
|
||||
self.retry_after.remove(&channel_id);
|
||||
self.retry_counts.remove(&channel_id);
|
||||
self.cancelled_batches.remove(&channel_id);
|
||||
self.cancel_reasons.remove(&channel_id);
|
||||
self.withheld_native_steer.remove(&channel_id);
|
||||
self.sent_native_steers.remove(&channel_id);
|
||||
// If the prompt already completed behind a sent-steer acknowledgement
|
||||
// fence, removing the reservation settles that fence immediately.
|
||||
@@ -4280,6 +4287,25 @@ mod tests {
|
||||
assert_eq!(pending_count(&q), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn drain_channel_returns_visible_reaction_targets_for_queued_and_reserved_edits() {
|
||||
let mut q = EventQueue::new(DedupMode::Queue);
|
||||
let ch = Uuid::new_v4();
|
||||
let queued_target = "ab".repeat(32);
|
||||
let reserved_target = "cd".repeat(32);
|
||||
|
||||
q.push(make_edit_queued_created_at(ch, &queued_target, 100));
|
||||
let reserved = make_edit_queued_created_at(ch, &reserved_target, 101);
|
||||
let reserved_id = reserved.event.id.to_hex();
|
||||
q.push(reserved);
|
||||
assert!(q.mark_native_steer_pending(ch, &reserved_id));
|
||||
|
||||
let drained = q.drain_channel(ch);
|
||||
assert_eq!(drained, vec![queued_target, reserved_target]);
|
||||
assert!(!q.has_native_steer_reservations(ch));
|
||||
assert_eq!(pending_count(&q), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_drain_channel_does_not_affect_other_channels() {
|
||||
let mut q = EventQueue::new(DedupMode::Queue);
|
||||
|
||||
Reference in New Issue
Block a user