diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 1a6661e22..406e4af36 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -2544,8 +2544,10 @@ async fn tokio_main() -> Result<()> { buzz_event.channel_id, event_for_steer, prompt_tag_for_steer, + &ctx, &steer_ack_tx, - ); + ) + .await; if !native_attempted { signal_in_flight_task( &mut pool, @@ -3158,40 +3160,21 @@ fn signal_in_flight_task( /// /// The withheld event is NOT released here on `false` because no withhold /// was established: `mark_native_steer_pending` only runs on `Ok(())`. -fn try_native_steer( +async fn try_native_steer( pool: &mut AgentPool, queue: &mut EventQueue, channel_id: uuid::Uuid, event: nostr::Event, prompt_tag: String, + ctx: &PromptContext, steer_ack_tx: &mpsc::UnboundedSender, ) -> bool { - // Build the steer body: framing strings come from - // `queue::native_steer_framing()` (Eva's drift-proof requirement — - // native and cancel+merge fallback share these so the agent gets the - // same orientation regardless of transport). The single event block - // is rendered by `queue::format_event_block`, the same function - // `queue::format_prompt` uses internally for `[Buzz event: …]` - // sections, so the rendering also cannot drift. - // - // Passing `None` for `channel_info` / `profile_lookup` is intentional: - // native steer is a *delta* into a live turn — the agent already saw - // channel context and the actor's profile in the original prompt, - // duplicating it here would defeat the point of non-cancelling - // steering (which is to inject only what's new). - let (header, closing) = queue::native_steer_framing(); let event_id_hex = event.id.to_hex(); - let be = queue::BatchEvent { - event, - prompt_tag: prompt_tag.clone(), - received_at: std::time::Instant::now(), - }; - let event_block = queue::format_event_block(channel_id, None, &be, None); - let body = format!("{header}\n\n[Buzz event: {prompt_tag}]\n{event_block}\n\n{closing}"); + let prompt_blocks = pool::format_native_steer_prompt(channel_id, event, prompt_tag, ctx).await; let (ack_tx, ack_rx) = tokio::sync::oneshot::channel::(); let request = pool::SteerRequest { - prompt_blocks: vec![body], + prompt_blocks, ack_tx, }; diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index e56855148..f9c393fd8 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -2888,7 +2888,7 @@ fn conversation_context_delta( /// Resolve a kind:40003 edit through its original event for reply routing. /// Failure deliberately falls back to the edit target id in `format_prompt`. -async fn resolve_edit_routing( +pub(crate) async fn resolve_edit_routing( event: &nostr::Event, rest: &RestClient, ) -> Option { @@ -2927,6 +2927,73 @@ async fn resolve_edit_routing( }) } +/// Resolve an edit-aware native-steer delta through the same routing and +/// context boundary as ordinary batch dispatch. This keeps successful native +/// steering from exposing the invisible kind:40003 event as a reply anchor. +pub(crate) async fn format_native_steer_prompt( + channel_id: Uuid, + event: nostr::Event, + prompt_tag: String, + ctx: &PromptContext, +) -> Vec { + let batch = FlushBatch { + channel_id, + events: vec![crate::queue::BatchEvent { + event, + prompt_tag, + received_at: std::time::Instant::now(), + }], + cancelled_events: Vec::new(), + cancel_reason: None, + }; + if crate::queue::edit_target_id(&batch.events[0].event).is_none() { + let (header, closing) = crate::queue::native_steer_framing(); + let event = &batch.events[0]; + return vec![format!( + "{header}\n\n[Buzz event: {}]\n{}\n\n{closing}", + event.prompt_tag, + crate::queue::format_event_block(channel_id, None, event, None) + )]; + } + let channel_info = ctx.channel_info.resolve(channel_id).await; + let resolved_edit = resolve_edit_routing(&batch.events[0].event, &ctx.rest_client).await; + let conversation_context = if ctx.context_message_limit > 0 { + match resolved_edit + .as_ref() + .and_then(|edit| edit.target_thread_tags.root_event_id.as_deref()) + { + Some(root_id) => { + fetch_thread_context( + channel_id, + root_id, + ctx.context_message_limit, + ctx.agent_keys.public_key(), + &ctx.rest_client, + ) + .await + } + None => fetch_conversation_context(&batch, &channel_info, ctx).await, + } + } else { + None + }; + let profile_lookup = + fetch_prompt_profile_lookup(&batch, conversation_context.as_ref(), &ctx.rest_client).await; + + crate::queue::format_native_steer_prompt( + &batch, + &crate::queue::FormatPromptArgs { + channel_info: channel_info.as_ref(), + conversation_context: conversation_context.as_ref(), + profile_lookup: profile_lookup.as_ref(), + resolved_edit: resolved_edit.as_ref(), + has_system_prompt_support: true, + standing_context_sent: true, + ..Default::default() + }, + ) +} + /// Fetch conversation context (thread or DM) for a batch before prompting. /// /// Returns `None` if: diff --git a/crates/buzz-acp/src/queue.rs b/crates/buzz-acp/src/queue.rs index a7b69883e..459351f5c 100644 --- a/crates/buzz-acp/src/queue.rs +++ b/crates/buzz-acp/src/queue.rs @@ -1729,6 +1729,21 @@ pub(crate) fn native_steer_framing() -> (&'static str, &'static str) { (framing.new_header_single, framing.closing_note) } +/// Format the delta delivered through native steering. +/// +/// Routing and context use the same formatter as ordinary dispatch, while the +/// native framing remains a concise "weave this into the live turn" envelope. +pub(crate) fn format_native_steer_prompt( + batch: &FlushBatch, + args: &FormatPromptArgs<'_>, +) -> Vec { + let (header, closing) = native_steer_framing(); + let mut sections = format_prompt(batch, args); + sections.insert(0, header.to_string()); + sections.push(closing.to_string()); + sections +} + #[cfg(test)] mod tests { use super::*; @@ -5056,6 +5071,46 @@ mod tests { assert!(!prompt.contains(&format!("--reply-to {edit_id}"))); } + #[test] + fn native_steer_edit_uses_original_thread_anchor() { + let original_id = "66".repeat(32); + let root_id = "77".repeat(32); + let batch = one_event_batch(edit_event(&original_id)); + let edit_id = batch.events[0].event.id.to_hex(); + let prompt = format_native_steer_prompt( + &batch, + &FormatPromptArgs { + resolved_edit: Some(&ResolvedEdit { + target_event_id: original_id, + target_thread_tags: ThreadTags { + root_event_id: Some(root_id.clone()), + parent_event_id: Some("88".repeat(32)), + mentioned_pubkeys: vec![], + }, + }), + ..Default::default() + }, + ) + .join("\n"); + + assert!(prompt.contains(&format!("--reply-to {root_id}"))); + assert!(!prompt.contains(&format!("--reply-to {edit_id}"))); + let (header, closing) = native_steer_framing(); + assert!(prompt.contains(header)); + assert!(prompt.contains(closing)); + } + + #[test] + fn native_steer_edit_fetch_failure_anchors_target_never_auxiliary_event() { + let original_id = "99".repeat(32); + let batch = one_event_batch(edit_event(&original_id)); + let edit_id = batch.events[0].event.id.to_hex(); + let prompt = format_native_steer_prompt(&batch, &FormatPromptArgs::default()).join("\n"); + + assert!(prompt.contains(&format!("--reply-to {original_id}"))); + assert!(!prompt.contains(&format!("--reply-to {edit_id}"))); + } + // ── F2 case 2.5: steer renewal is monotonic across repeated steers ─────── /// Calling `extend_in_flight_deadline` twice with the same `max_turn_secs`